use std::sync::Arc;
use chrono::{DateTime, Utc};
use sqlx::PgPool;
use uuid::Uuid;
use crate::application::service::digest_error::DigestError;
use crate::application::service::digest_render_service::{active_recipients, DigestRenderService};
use crate::application::service::digest_write_service::{
add_months, DigestPeriodicity, DigestRow,
};
use crate::application::service::engagement_port::{RecipientContextPort, RecipientContextSlot};
impl DigestPeriodicity {
pub fn login_window_since(&self, now: DateTime<Utc>) -> DateTime<Utc> {
match self {
DigestPeriodicity::Daily => now - chrono::Duration::days(2),
DigestPeriodicity::Weekly => now - chrono::Duration::days(7),
DigestPeriodicity::Monthly => add_months(now.date_naive(), -1)
.and_hms_opt(0, 0, 0)
.expect("valid")
.and_utc(),
DigestPeriodicity::Quarterly => add_months(now.date_naive(), -3)
.and_hms_opt(0, 0, 0)
.expect("valid")
.and_utc(),
}
}
}
#[derive(Debug, Default, Clone)]
pub struct SweepOutcome {
pub claimed: usize,
pub sent_mails: usize,
pub degraded: Vec<Uuid>,
pub mail_delivery_failures: Vec<Uuid>,
pub isolated_failures: Vec<Uuid>,
}
pub struct DigestCronService {
pool: PgPool,
render: Arc<DigestRenderService>,
port: RecipientContextSlot,
}
impl DigestCronService {
pub fn new(pool: PgPool, render: Arc<DigestRenderService>, port: RecipientContextSlot) -> Self {
Self { pool, render, port }
}
pub async fn run_daily_pull(&self, now: DateTime<Utc>) -> Result<SweepOutcome, DigestError> {
let today = now.date_naive();
let mut outcome = SweepOutcome::default();
let mut processed: Vec<Uuid> = Vec::new();
loop {
let mut tx = self.pool.begin().await?;
let claimed = sqlx::query_as::<_, (Uuid, String, String, Option<chrono::NaiveDate>, String)>(
r#"SELECT id, name, periodicity::text, next_run_date, state::text
FROM digest.digest_digests
WHERE state = 'activated'
AND next_run_date IS NOT NULL
AND next_run_date <= $1
AND NOT (id = ANY($2))
ORDER BY next_run_date, id
FOR UPDATE SKIP LOCKED
LIMIT 1"#,
)
.bind(today)
.bind(&processed)
.fetch_optional(&mut *tx)
.await?;
let Some((id, name, p, next_run_date, state)) = claimed else {
break;
};
let digest = DigestRow {
id,
name,
periodicity: DigestPeriodicity::parse(&p).unwrap_or(DigestPeriodicity::Daily),
next_run_date,
state,
};
processed.push(digest.id);
outcome.claimed += 1;
match self.process_digest(&mut tx, &digest, now, &mut outcome).await {
Ok(()) => { tx.commit().await?; }
Err(e) => {
tx.rollback().await?;
tracing::error!(
target: "digest::audit",
event = "digest_send_isolated",
digest_id = %digest.id,
error = %e,
"digest send failed; digest left due (no advance), sweep continues"
);
outcome.isolated_failures.push(digest.id);
}
}
}
Ok(outcome)
}
async fn process_digest(
&self,
tx: &mut sqlx::PgTransaction<'_>,
digest: &DigestRow,
now: DateTime<Utc>,
outcome: &mut SweepOutcome,
) -> Result<(), DigestError> {
let recipients = active_recipients(&self.pool, digest.id).await?;
let mut cadence = digest.periodicity;
let user_ids: Vec<Uuid> = recipients.iter().map(|r| r.user_id).collect();
let since = digest.periodicity.login_window_since(now);
match self.port.any_logged_in_since(&user_ids, since).await {
Ok(false) => {
if let Some(next) = cadence.next_rung() {
sqlx::query(
"UPDATE digest.digest_digests SET periodicity = $2::digest_periodicity WHERE id = $1",
)
.bind(digest.id)
.bind(next.as_str())
.execute(&mut **tx)
.await?;
tracing::info!(
target: "digest::audit",
event = "digest_slowdown_degraded",
digest_id = %digest.id,
from = cadence.as_str(),
to = next.as_str(),
recipients = user_ids.len()
);
cadence = next;
outcome.degraded.push(digest.id);
}
}
Ok(true) => { }
Err(e) => {
tracing::error!(
target: "digest::audit",
event = "slowdown_signal_unavailable",
digest_id = %digest.id,
error = %e,
"slowdown signal unavailable through the recipient port; cadence held"
);
}
}
let mut mail_failure = false;
let mut digest_failure: Option<String> = None;
for recipient in &recipients {
match self.render.send_to_recipient(digest, recipient, now).await {
Ok((_plan, _mail_id)) => {
outcome.sent_mails += 1;
}
Err(DigestError::MailDelivery(why)) => {
tracing::warn!(
digest_id = %digest.id,
user_id = %recipient.user_id,
reason = %why,
"mail delivery refused for recipient; digest will retry whole next day (no advance)"
);
mail_failure = true;
}
Err(e) => {
tracing::error!(
digest_id = %digest.id,
user_id = %recipient.user_id,
error = %e,
"unexpected recipient send failure (isolated)"
);
digest_failure = Some(e.to_string());
}
}
}
if mail_failure || digest_failure.is_some() {
if mail_failure {
outcome.mail_delivery_failures.push(digest.id);
} else {
outcome.isolated_failures.push(digest.id);
}
return Ok(());
}
let advanced_to = cadence.advance(now.date_naive());
let n = sqlx::query("UPDATE digest.digest_digests SET next_run_date = $2 WHERE id = $1")
.bind(digest.id)
.bind(advanced_to)
.execute(&mut **tx)
.await?
.rows_affected();
if n == 0 {
return Err(DigestError::SendFailed(format!(
"digest {} vanished mid-batch (advance matched 0 rows)",
digest.id
)));
}
Ok(())
}
pub async fn send_now(
&self,
digest_id: Uuid,
now: DateTime<Utc>,
) -> Result<(usize, usize, usize), DigestError> {
let Some(digest) = (|| async {
let row = sqlx::query_as::<_, (Uuid, String, String, Option<chrono::NaiveDate>, String)>(
r#"SELECT id, name, periodicity::text, next_run_date, state::text
FROM digest.digest_digests WHERE id = $1"#,
)
.bind(digest_id)
.fetch_optional(&self.pool)
.await?;
Ok::<_, DigestError>(row.map(|(id, name, p, next_run_date, state)| DigestRow {
id,
name,
periodicity: DigestPeriodicity::parse(&p).unwrap_or(DigestPeriodicity::Daily),
next_run_date,
state,
}))
})()
.await?
else {
return Err(DigestError::NotFound(digest_id));
};
let recipients = active_recipients(&self.pool, digest.id).await?;
let (mut sent, mut mail_fail, mut unexpected) = (0, 0, 0);
for recipient in &recipients {
match self.render.send_to_recipient(&digest, recipient, now).await {
Ok(_) => sent += 1,
Err(DigestError::MailDelivery(why)) => {
tracing::warn!(digest_id = %digest.id, reason = %why, "manual send: mail delivery refused");
mail_fail += 1;
}
Err(e) => {
tracing::error!(digest_id = %digest.id, error = %e, "manual send: unexpected failure");
unexpected += 1;
}
}
}
Ok((sent, mail_fail, unexpected))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn ladder_windows_follow_the_rule_table() {
let now = DateTime::parse_from_rfc3339("2026-08-31T02:41:00Z").unwrap().with_timezone(&Utc);
let d = DigestPeriodicity::Daily.login_window_since(now);
assert_eq!((now - d).num_days(), 2);
let w = DigestPeriodicity::Weekly.login_window_since(now);
assert_eq!((now - w).num_days(), 7);
let m = DigestPeriodicity::Monthly.login_window_since(now);
assert_eq!(m.date_naive(), chrono::NaiveDate::from_ymd_opt(2026, 7, 31).unwrap());
let q = DigestPeriodicity::Quarterly.login_window_since(now);
assert_eq!(q.date_naive(), chrono::NaiveDate::from_ymd_opt(2026, 5, 31).unwrap());
}
}