systemprompt_logging/repository/analytics/
mod.rs1use chrono::Utc;
12use serde_json::Value;
13use sqlx::PgPool;
14use std::sync::Arc;
15use systemprompt_database::DbPool;
16use systemprompt_identifiers::{AgentId, ContextId, SessionId, TaskId, UserId};
17
18use crate::models::LoggingError;
19
20#[derive(Debug, Clone)]
21pub struct AnalyticsRepository {
22 write_pool: Arc<PgPool>,
23}
24
25impl AnalyticsRepository {
26 pub fn new(db: &DbPool) -> Result<Self, LoggingError> {
27 let write_pool = db.write_pool_arc()?;
28 Ok(Self { write_pool })
29 }
30
31 pub async fn log_event(&self, event: &AnalyticsEvent) -> Result<i64, LoggingError> {
32 let result = execute_insert(&self.write_pool, event).await?;
33 Ok(i64::try_from(result).unwrap_or(i64::MAX))
34 }
35}
36
37async fn execute_insert(pool: &PgPool, event: &AnalyticsEvent) -> Result<u64, LoggingError> {
38 let params = EventParams::from(event);
39 run_insert_query(pool, params).await
40}
41
42async fn run_insert_query(pool: &PgPool, p: EventParams<'_>) -> Result<u64, LoggingError> {
43 let result = sqlx::query!(
44 r"
45 INSERT INTO analytics_events
46 (user_id, session_id, context_id, event_type, event_category, severity,
47 endpoint, error_code, response_time_ms, agent_id, task_id, message, metadata, timestamp)
48 VALUES
49 ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14)
50 ",
51 p.user_id,
52 p.session_id,
53 p.context_id,
54 p.event_type,
55 p.event_category,
56 p.severity,
57 p.endpoint,
58 p.error_code,
59 p.response_time_ms,
60 p.agent_id,
61 p.task_id,
62 p.message,
63 p.metadata,
64 p.timestamp
65 )
66 .execute(pool)
67 .await?;
68 Ok(result.rows_affected())
69}
70
71struct EventParams<'a> {
72 user_id: &'a str,
73 session_id: &'a str,
74 context_id: &'a str,
75 event_type: &'a str,
76 event_category: &'a str,
77 severity: &'a str,
78 agent_id: Option<&'a str>,
79 task_id: Option<&'a str>,
80 endpoint: Option<&'a str>,
81 message: Option<&'a str>,
82 error_code: Option<i32>,
83 response_time_ms: Option<i32>,
84 metadata: String,
85 timestamp: chrono::DateTime<Utc>,
86}
87
88impl<'a> From<&'a AnalyticsEvent> for EventParams<'a> {
89 fn from(event: &'a AnalyticsEvent) -> Self {
90 Self {
91 user_id: event.user_id.as_str(),
92 session_id: event.session_id.as_str(),
93 context_id: event.context_id.as_str(),
94 event_type: &event.event_type,
95 event_category: &event.event_category,
96 severity: &event.severity,
97 agent_id: event.agent_id.as_ref().map(AgentId::as_str),
98 task_id: event.task_id.as_ref().map(TaskId::as_str),
99 endpoint: event.endpoint.as_deref(),
100 message: event.message.as_deref(),
101 error_code: event.error_code,
102 response_time_ms: event.response_time_ms,
103 metadata: event.metadata.to_string(),
104 timestamp: Utc::now(),
105 }
106 }
107}
108
109#[derive(Debug, Clone)]
110pub struct AnalyticsEvent {
111 pub user_id: UserId,
112 pub session_id: SessionId,
113 pub context_id: ContextId,
114 pub event_type: String,
115 pub event_category: String,
116 pub severity: String,
117 pub endpoint: Option<String>,
118 pub error_code: Option<i32>,
119 pub response_time_ms: Option<i32>,
120 pub agent_id: Option<AgentId>,
121 pub task_id: Option<TaskId>,
122 pub message: Option<String>,
123 pub metadata: Value,
124}