use std::collections::BTreeMap;
use std::sync::Arc;
use chrono::{DateTime, Utc};
use sqlx::PgPool;
use uuid::Uuid;
pub const KPI_NAME_PREFIX: &str = "kpi_";
pub const KPI_CONNECTED_USERS: &str = "kpi_res_users_connected";
pub const KPI_MESSAGES_SENT: &str = "kpi_mail_message_total";
#[derive(Debug, Clone)]
pub struct KpiQueryContext {
pub recipient_user_id: Uuid,
pub company_id: Option<Uuid>,
pub window_start: DateTime<Utc>,
pub window_end: DateTime<Utc>,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct KpiValue {
pub count: i64,
pub amount: Option<f64>,
}
impl KpiValue {
pub fn zero() -> Self {
Self::default()
}
pub fn count(n: i64) -> Self {
Self { count: n, amount: None }
}
pub fn display(&self) -> f64 {
self.amount.unwrap_or(self.count as f64)
}
}
#[derive(Debug, thiserror::Error)]
pub enum KpiError {
#[error("kpi compute failed: {0}")]
Compute(#[source] sqlx::Error),
#[error("kpi source unavailable: {0}")]
Unavailable(String),
}
#[async_trait::async_trait]
pub trait KpiComputer: Send + Sync {
async fn compute(&self, pool: &PgPool, ctx: &KpiQueryContext) -> Result<KpiValue, KpiError>;
}
pub struct FnComputer<F>(pub F);
#[async_trait::async_trait]
impl<F, Fut> KpiComputer for FnComputer<F>
where
F: Fn(&PgPool, KpiQueryContext) -> Fut + Send + Sync,
Fut: std::future::Future<Output = Result<KpiValue, KpiError>> + Send,
{
async fn compute(&self, pool: &PgPool, ctx: &KpiQueryContext) -> Result<KpiValue, KpiError> {
(self.0)(pool, ctx.clone()).await
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum KpiFenceDeclaration {
SharedData,
CompanyData,
PinnedCompanyData(Uuid),
}
impl KpiFenceDeclaration {
pub fn renders_for(&self, recipient_company: Option<Uuid>) -> bool {
match self {
KpiFenceDeclaration::SharedData => true,
KpiFenceDeclaration::CompanyData => recipient_company.is_some(),
KpiFenceDeclaration::PinnedCompanyData(pinned) => {
recipient_company == Some(*pinned)
}
}
}
}
#[derive(Clone)]
pub struct KpiDefinition {
pub name: String,
pub label: String,
pub fence: KpiFenceDeclaration,
pub computer: Arc<dyn KpiComputer>,
}
impl std::fmt::Debug for KpiDefinition {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("KpiDefinition")
.field("name", &self.name)
.field("label", &self.label)
.field("fence", &self.fence)
.finish_non_exhaustive()
}
}
#[derive(Debug, Default)]
pub struct KpiRegistry {
entries: std::sync::RwLock<BTreeMap<String, KpiDefinition>>,
}
#[derive(Debug, thiserror::Error, PartialEq)]
pub enum KpiRegistrationError {
#[error("KPI name '{0}' must start with '{1}'")]
BadPrefix(String, String),
}
impl KpiRegistry {
pub fn new() -> Self {
Self::default()
}
pub fn register(
&self,
def: KpiDefinition,
) -> Result<(), KpiRegistrationError> {
if !def.name.starts_with(KPI_NAME_PREFIX) {
return Err(KpiRegistrationError::BadPrefix(
def.name,
KPI_NAME_PREFIX.to_string(),
));
}
let mut guard = self.entries.write().unwrap_or_else(|e| e.into_inner());
if guard.contains_key(&def.name) {
panic!(
"duplicate KPI registration: '{}' is already registered — \
a metric name is an identity, not a slot (R-DG1)",
def.name
);
}
guard.insert(def.name.clone(), def);
Ok(())
}
pub fn get(&self, name: &str) -> Option<KpiDefinition> {
self.entries
.read()
.unwrap_or_else(|e| e.into_inner())
.get(name)
.cloned()
}
pub fn names(&self) -> Vec<String> {
self.entries
.read()
.unwrap_or_else(|e| e.into_inner())
.keys()
.cloned()
.collect()
}
pub fn len(&self) -> usize {
self.entries
.read()
.unwrap_or_else(|e| e.into_inner())
.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}
#[cfg(test)]
mod tests {
use super::*;
struct StubComputer;
#[async_trait::async_trait]
impl KpiComputer for StubComputer {
async fn compute(&self, _pool: &PgPool, _ctx: &KpiQueryContext) -> Result<KpiValue, KpiError> {
Ok(KpiValue::count(7))
}
}
fn def(name: &str, fence: KpiFenceDeclaration) -> KpiDefinition {
KpiDefinition {
name: name.into(),
label: name.into(),
fence,
computer: Arc::new(StubComputer),
}
}
#[test]
fn prefix_is_enforced() {
let r = KpiRegistry::new();
assert_eq!(
r.register(def("crm_leads", KpiFenceDeclaration::SharedData)),
Err(KpiRegistrationError::BadPrefix(
"crm_leads".into(),
"kpi_".into()
))
);
assert!(r.register(def("kpi_crm_leads", KpiFenceDeclaration::SharedData)).is_ok());
}
#[test]
#[should_panic(expected = "duplicate KPI registration")]
fn duplicate_registration_panics() {
let r = KpiRegistry::new();
r.register(def("kpi_crm_leads", KpiFenceDeclaration::SharedData)).unwrap();
let _ = r.register(def("kpi_crm_leads", KpiFenceDeclaration::CompanyData));
}
#[test]
fn shared_renders_for_everyone() {
assert!(KpiFenceDeclaration::SharedData.renders_for(None));
assert!(KpiFenceDeclaration::SharedData.renders_for(Some(Uuid::new_v4())));
}
#[test]
fn company_data_renders_only_for_recipients_with_a_company() {
let a = Uuid::new_v4();
let f = KpiFenceDeclaration::CompanyData;
assert!(f.renders_for(Some(a)));
assert!(!f.renders_for(None));
}
#[test]
fn pinned_company_renders_only_for_that_company() {
let pinned = Uuid::new_v4();
let f = KpiFenceDeclaration::PinnedCompanyData(pinned);
assert!(f.renders_for(Some(pinned)));
assert!(!f.renders_for(Some(Uuid::new_v4())));
assert!(!f.renders_for(None));
}
}