use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
pub struct AnalyticsEvent {
pub event_id: String,
pub property_id: String,
pub tenant_id: String,
pub event_name: String,
pub event_source: String,
pub dimensions: BTreeMap<String, String>,
pub metrics: BTreeMap<String, f64>,
pub occurred_at_unix: u64,
pub received_at_unix: u64,
}
impl AnalyticsEvent {
pub fn is_valid(&self) -> bool {
!self.event_id.is_empty() && !self.event_name.is_empty()
}
pub fn validate(&self) -> Option<String> {
if self.event_id.is_empty() {
return Some("event_id is required".to_string());
}
if self.event_name.is_empty() {
return Some("event_name is required".to_string());
}
None
}
}
pub const DDL_ANALYTICS_EVENTS_DAILY: &str = r#"
CREATE SCHEMA IF NOT EXISTS analytics;
CREATE TABLE IF NOT EXISTS analytics.events_daily (
event_id TEXT NOT NULL,
property_id TEXT NOT NULL DEFAULT '',
tenant_id TEXT NOT NULL DEFAULT '',
event_name TEXT NOT NULL,
event_source TEXT NOT NULL DEFAULT 'internal',
dimensions JSONB NULL,
metrics JSONB NULL,
occurred_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
received_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
CONSTRAINT pk_events_daily PRIMARY KEY (event_id, occurred_at)
) PARTITION BY RANGE (occurred_at);
-- Default partition for any day not yet explicitly partitioned.
CREATE TABLE IF NOT EXISTS analytics.events_daily_default
PARTITION OF analytics.events_daily DEFAULT;
CREATE INDEX IF NOT EXISTS idx_events_daily_event_name
ON analytics.events_daily (event_name, occurred_at DESC);
CREATE INDEX IF NOT EXISTS idx_events_daily_tenant
ON analytics.events_daily (tenant_id, occurred_at DESC);
"#;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
pub struct SignalConfig {
pub query_timeout_secs: u64,
pub batch_size: usize,
pub fallback_to_primary: bool,
}
impl SignalConfig {
pub fn effective_query_timeout_secs(&self) -> u64 {
if self.query_timeout_secs > 0 {
self.query_timeout_secs
} else {
30
}
}
pub fn effective_batch_size(&self) -> usize {
if self.batch_size > 0 {
self.batch_size
} else {
500
}
}
}
pub mod common_events {
pub const DOCUMENT_SUBMITTED: &str = "document.submitted.v1";
pub const DOCUMENT_PROCESSED: &str = "document.processed.v1";
pub const DOCUMENT_FAILED: &str = "document.failed.v1";
pub const FIELD_CORRECTED: &str = "field.corrected.v1";
pub const TEMPLATE_PUBLISHED: &str = "template.published.v1";
pub const LEARNING_CYCLE_COMPLETED: &str = "learning.cycle_completed.v1";
pub const MIGRATION_COMPLETED: &str = "udb.migration.completed.v1";
pub const MIGRATION_FAILED: &str = "udb.migration.failed.v1";
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn analytics_event_validate_empty_id() {
let evt = AnalyticsEvent::default();
assert!(evt.validate().is_some());
}
#[test]
fn analytics_event_validate_valid() {
let evt = AnalyticsEvent {
event_id: "evt-001".to_string(),
event_name: common_events::DOCUMENT_PROCESSED.to_string(),
..Default::default()
};
assert!(evt.validate().is_none());
assert!(evt.is_valid());
}
#[test]
fn signal_config_defaults() {
let cfg = SignalConfig::default();
assert_eq!(cfg.effective_query_timeout_secs(), 30);
assert_eq!(cfg.effective_batch_size(), 500);
}
#[test]
fn ddl_analytics_events_daily_nonempty() {
assert!(!DDL_ANALYTICS_EVENTS_DAILY.trim().is_empty());
assert!(DDL_ANALYTICS_EVENTS_DAILY.contains("analytics.events_daily"));
}
#[test]
fn common_event_names_are_versioned() {
for name in [
common_events::DOCUMENT_SUBMITTED,
common_events::DOCUMENT_PROCESSED,
common_events::DOCUMENT_FAILED,
common_events::FIELD_CORRECTED,
common_events::TEMPLATE_PUBLISHED,
common_events::MIGRATION_COMPLETED,
common_events::MIGRATION_FAILED,
] {
assert!(
name.ends_with(".v1"),
"event name '{name}' must end with .v1"
);
}
}
}