use backbone_orm::org_scope;
use chrono::{Datelike, NaiveDate, Utc};
use sqlx::{PgPool, Row};
use uuid::Uuid;
use crate::application::service::digest_error::DigestError;
use crate::application::service::kpi_registry::KpiRegistry;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum DigestPeriodicity {
Daily,
Weekly,
Monthly,
Quarterly,
}
impl DigestPeriodicity {
pub fn parse(value: &str) -> Option<Self> {
match value {
"daily" => Some(DigestPeriodicity::Daily),
"weekly" => Some(DigestPeriodicity::Weekly),
"monthly" => Some(DigestPeriodicity::Monthly),
"quarterly" => Some(DigestPeriodicity::Quarterly),
_ => None,
}
}
pub fn as_str(&self) -> &'static str {
match self {
DigestPeriodicity::Daily => "daily",
DigestPeriodicity::Weekly => "weekly",
DigestPeriodicity::Monthly => "monthly",
DigestPeriodicity::Quarterly => "quarterly",
}
}
pub fn advance(&self, from: NaiveDate) -> NaiveDate {
match self {
DigestPeriodicity::Daily => from + chrono::Duration::days(1),
DigestPeriodicity::Weekly => from + chrono::Duration::weeks(1),
DigestPeriodicity::Monthly => add_months(from, 1),
DigestPeriodicity::Quarterly => add_months(from, 3),
}
}
pub fn next_rung(&self) -> Option<DigestPeriodicity> {
match self {
DigestPeriodicity::Daily => Some(DigestPeriodicity::Weekly),
DigestPeriodicity::Weekly => Some(DigestPeriodicity::Monthly),
DigestPeriodicity::Monthly => Some(DigestPeriodicity::Quarterly),
DigestPeriodicity::Quarterly => None,
}
}
}
pub fn add_months(from: NaiveDate, months: i32) -> NaiveDate {
let total = from.year() * 12 + (from.month0() as i32) + months;
let year = total.div_euclid(12);
let month0 = total.rem_euclid(12) as u32;
let day = from.day().min(days_in_month(year, month0 + 1));
NaiveDate::from_ymd_opt(year, month0 + 1, day).expect("valid ymd by construction")
}
fn days_in_month(year: i32, month: u32) -> u32 {
let (y, m) = if month == 12 { (year + 1, 1) } else { (year, month + 1) };
let first_of_next = NaiveDate::from_ymd_opt(y, m, 1).expect("valid first of month");
(first_of_next - chrono::Duration::days(1)).day()
}
#[derive(Debug, Clone)]
pub struct DigestRow {
pub id: Uuid,
pub name: String,
pub periodicity: DigestPeriodicity,
pub next_run_date: Option<NaiveDate>,
pub state: String,
}
impl DigestRow {
pub fn is_activated(&self) -> bool {
self.state == "activated"
}
}
pub struct DigestWriteService {
pool: PgPool,
registry: std::sync::Arc<KpiRegistry>,
}
impl DigestWriteService {
pub fn new(pool: PgPool, registry: std::sync::Arc<KpiRegistry>) -> Self {
Self { pool, registry }
}
pub async fn create_digest(
&self,
name: &str,
periodicity: DigestPeriodicity,
today: NaiveDate,
) -> Result<Uuid, DigestError> {
let id = Uuid::new_v4();
org_scope::execute_scoped(
&self.pool,
sqlx::query(
r#"INSERT INTO digest.digest_digests
(id, name, periodicity, next_run_date, state)
VALUES ($1, $2, $3::digest_periodicity, $4, 'activated')"#,
)
.bind(id)
.bind(name)
.bind(periodicity.as_str())
.bind(periodicity.advance(today)),
)
.await?;
Ok(id)
}
pub async fn set_periodicity(
&self,
digest_id: Uuid,
value: &str,
today: NaiveDate,
) -> Result<(), DigestError> {
let Some(p) = DigestPeriodicity::parse(value) else {
return Err(DigestError::Invalid(format!(
"periodicity '{value}' refused — the whitelist is daily|weekly|monthly|quarterly (error code: periodicity_value_refused)"
)));
};
let n = org_scope::execute_scoped(
&self.pool,
sqlx::query(
r#"UPDATE digest.digest_digests
SET periodicity = $2::digest_periodicity,
next_run_date = $3
WHERE id = $1"#,
)
.bind(digest_id)
.bind(p.as_str())
.bind(p.advance(today)),
)
.await?
.rows_affected();
if n == 0 {
return Err(DigestError::NotFound(digest_id));
}
Ok(())
}
pub async fn enable_kpi(&self, digest_id: Uuid, kpi_key: &str) -> Result<(), DigestError> {
if self.registry.get(kpi_key).is_none() {
return Err(DigestError::Invalid(format!(
"KPI '{kpi_key}' is not in the composed registry (error code: kpi_not_registered)"
)));
}
org_scope::execute_scoped(
&self.pool,
sqlx::query(
r#"INSERT INTO digest.digest_digest_kpis (digest_id, kpi_key)
VALUES ($1, $2)
ON CONFLICT (digest_id, kpi_key)
DO NOTHING"#,
)
.bind(digest_id)
.bind(kpi_key),
)
.await?;
Ok(())
}
pub async fn disable_kpi(&self, digest_id: Uuid, kpi_key: &str) -> Result<(), DigestError> {
org_scope::execute_scoped(
&self.pool,
sqlx::query(
r#"DELETE FROM digest.digest_digest_kpis
WHERE digest_id = $1 AND kpi_key = $2"#,
)
.bind(digest_id)
.bind(kpi_key),
)
.await?;
Ok(())
}
pub async fn subscribe_user(
&self,
digest_id: Uuid,
user_id: Uuid,
extra_metadata: Option<serde_json::Value>,
) -> Result<bool, DigestError> {
let empty = serde_json::json!({});
let meta = extra_metadata.unwrap_or(empty);
let inserted = sqlx::query(
r#"INSERT INTO digest.digest_subscriptions (digest_id, user_id, state, metadata)
VALUES ($1, $2, 'subscribed', $3::jsonb)
ON CONFLICT (digest_id, user_id)
DO UPDATE SET state = 'subscribed',
unsubscribed_at = NULL,
metadata = (digest.digest_subscriptions.metadata - 'deleted_at')
|| excluded.metadata
RETURNING (xmax = 0) AS inserted"#,
)
.bind(digest_id)
.bind(user_id)
.bind(meta)
.fetch_one(&self.pool)
.await?
.try_get::<bool, _>("inserted")
.unwrap_or(false);
Ok(inserted)
}
pub async fn set_state(&self, digest_id: Uuid, activated: bool) -> Result<(), DigestError> {
let value = if activated { "activated" } else { "deactivated" };
let n = org_scope::execute_scoped(
&self.pool,
sqlx::query(
r#"UPDATE digest.digest_digests
SET state = $2::digest_state
WHERE id = $1"#,
)
.bind(digest_id)
.bind(value),
)
.await?
.rows_affected();
if n == 0 {
return Err(DigestError::NotFound(digest_id));
}
Ok(())
}
pub async fn get_digest(&self, digest_id: Uuid) -> Result<Option<DigestRow>, DigestError> {
let row = sqlx::query_as::<_, (Uuid, String, String, Option<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(row.map(|(id, name, p, next_run_date, state)| DigestRow {
id,
name,
periodicity: DigestPeriodicity::parse(&p)
.unwrap_or(DigestPeriodicity::Daily),
next_run_date,
state,
}))
}
pub async fn stamp_subscription_metadata(
&self,
digest_id: Uuid,
user_id: Uuid,
patch: serde_json::Value,
) -> Result<(), DigestError> {
org_scope::execute_scoped(
&self.pool,
sqlx::query(
r#"UPDATE digest.digest_subscriptions
SET metadata = metadata || $3::jsonb
WHERE digest_id = $1 AND user_id = $2"#,
)
.bind(digest_id)
.bind(user_id)
.bind(patch),
)
.await?;
Ok(())
}
}
pub fn utc_today(now: chrono::DateTime<Utc>) -> NaiveDate {
now.date_naive()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn whitelist_is_exact() {
for ok in ["daily", "weekly", "monthly", "quarterly"] {
assert!(DigestPeriodicity::parse(ok).is_some(), "{ok} must parse");
}
for refused in ["Daily", "biweekly", "", "yearly", "hourly"] {
assert!(DigestPeriodicity::parse(refused).is_none(), "{refused} must be refused");
}
}
#[test]
fn ladder_order_is_monotonic_with_quarterly_floor() {
assert_eq!(DigestPeriodicity::Daily.next_rung(), Some(DigestPeriodicity::Weekly));
assert_eq!(DigestPeriodicity::Weekly.next_rung(), Some(DigestPeriodicity::Monthly));
assert_eq!(DigestPeriodicity::Monthly.next_rung(), Some(DigestPeriodicity::Quarterly));
assert_eq!(DigestPeriodicity::Quarterly.next_rung(), None);
}
#[test]
fn advance_uses_calendar_units() {
let d = NaiveDate::from_ymd_opt(2026, 1, 31).unwrap();
assert_eq!(DigestPeriodicity::Monthly.advance(d), NaiveDate::from_ymd_opt(2026, 2, 28).unwrap());
assert_eq!(DigestPeriodicity::Quarterly.advance(d), NaiveDate::from_ymd_opt(2026, 4, 30).unwrap());
assert_eq!(DigestPeriodicity::Daily.advance(d), NaiveDate::from_ymd_opt(2026, 2, 1).unwrap());
let dec31 = NaiveDate::from_ymd_opt(2026, 12, 31).unwrap();
assert_eq!(DigestPeriodicity::Monthly.advance(dec31), NaiveDate::from_ymd_opt(2027, 1, 31).unwrap());
}
}