use chrono::{DateTime, Utc};
use sqlx::PgPool;
use crate::application::service::engagement_port::{RecipientContextPort, RecipientContextSlot};
use crate::application::service::kpi_registry::{
KpiComputer, KpiError, KpiQueryContext, KpiValue, KPI_CONNECTED_USERS, KPI_MESSAGES_SENT,
};
pub struct ConnectedUsersKpi {
port: RecipientContextSlot,
}
impl ConnectedUsersKpi {
pub fn new(port: RecipientContextSlot) -> Self {
Self { port }
}
}
#[async_trait::async_trait]
impl KpiComputer for ConnectedUsersKpi {
async fn compute(&self, _pool: &PgPool, ctx: &KpiQueryContext) -> Result<KpiValue, KpiError> {
let n = self
.port
.count_connected(ctx.company_id, ctx.window_start, ctx.window_end)
.await
.map_err(|e| KpiError::Unavailable(e.to_string()))?;
Ok(KpiValue::count(n))
}
}
pub struct MessagesSentKpi;
const MESSAGES_SENT_SQL: &str = r#"
SELECT count(*) AS n
FROM messaging.mail_messages
WHERE message_type IN ('email', 'comment')
AND date >= $1
AND date < $2
AND (metadata->>'deleted_at') IS NULL
"#;
#[async_trait::async_trait]
impl KpiComputer for MessagesSentKpi {
async fn compute(&self, pool: &PgPool, ctx: &KpiQueryContext) -> Result<KpiValue, KpiError> {
let n: i64 = sqlx::query_scalar(MESSAGES_SENT_SQL)
.bind(ctx.window_start)
.bind(ctx.window_end)
.fetch_one(pool)
.await
.map_err(KpiError::Compute)?;
Ok(KpiValue::count(n))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct TimeWindow {
pub start: DateTime<Utc>,
pub end: DateTime<Utc>,
}
impl TimeWindow {
pub fn current(now: DateTime<Utc>, days: i64) -> Self {
Self { start: now - chrono::Duration::days(days), end: now }
}
pub fn previous(&self) -> Self {
let span = self.end - self.start;
Self { start: self.start - span, end: self.start }
}
}
pub const WINDOW_SPANS_DAYS: [i64; 3] = [1, 7, 30];
pub fn windows_at(now: DateTime<Utc>) -> [(TimeWindow, TimeWindow); 3] {
WINDOW_SPANS_DAYS
.map(|days| TimeWindow::current(now, days))
.map(|cur| (cur, cur.previous()))
}
pub fn margin(current: f64, previous: f64) -> Option<f64> {
if current != previous && current != 0.0 && previous != 0.0 {
Some(((current - previous) / previous * 100.0 * 100.0).round() / 100.0)
} else {
None
}
}
pub fn register_base_kpis(
registry: &crate::application::service::kpi_registry::KpiRegistry,
port: RecipientContextSlot,
) {
use crate::application::service::kpi_registry::{KpiDefinition, KpiFenceDeclaration};
let _ = registry.register(KpiDefinition {
name: KPI_CONNECTED_USERS.into(),
label: "Connected Users".into(),
fence: KpiFenceDeclaration::CompanyData,
computer: std::sync::Arc::new(ConnectedUsersKpi::new(port)),
});
let _ = registry.register(KpiDefinition {
name: KPI_MESSAGES_SENT.into(),
label: "Messages Sent".into(),
fence: KpiFenceDeclaration::SharedData,
computer: std::sync::Arc::new(MessagesSentKpi),
});
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn windows_are_utc_and_previous_shifts_exactly() {
let now = Utc::now();
let [(c1, p1), (c7, p7), (c30, p30)] = windows_at(now);
assert_eq!(c1.end, now);
assert_eq!((c1.end - c1.start).num_hours(), 24);
assert_eq!(p1.end, c1.start);
assert_eq!((p1.end - p1.start).num_hours(), 24);
assert_eq!((c7.end - c7.start).num_days(), 7);
assert_eq!(p7.end, c7.start);
assert_eq!((c30.end - c30.start).num_days(), 30);
assert_eq!(p30.end, c30.start);
}
#[test]
fn margin_only_when_changed_and_both_nonzero() {
assert_eq!(margin(10.0, 5.0), Some(100.0));
assert_eq!(margin(5.0, 10.0), Some(-50.0));
assert_eq!(margin(0.0, 5.0), None);
assert_eq!(margin(5.0, 0.0), None);
assert_eq!(margin(0.0, 0.0), None);
assert_eq!(margin(5.0, 5.0), None);
assert_eq!(margin(3.0, 9.0), Some(-66.67));
}
}