systemprompt_analytics/repository/
events.rs1use std::sync::Arc;
7
8use crate::Result;
9use sqlx::PgPool;
10use systemprompt_database::DbPool;
11use systemprompt_identifiers::{ContentId, SessionId, UserId};
12use systemprompt_traits::analytics_events::{AnalyticsEventRecord, DynAnalyticsEventStore};
13
14use crate::models::{AnalyticsEventCreated, AnalyticsEventType, CreateAnalyticsEventInput};
15
16#[derive(Clone, Debug)]
17pub struct AnalyticsEventsRepository {
18 pool: Arc<PgPool>,
19 event_sink: DynAnalyticsEventStore,
20}
21
22impl AnalyticsEventsRepository {
23 pub fn new(db: &DbPool, event_sink: DynAnalyticsEventStore) -> Result<Self> {
24 let pool = db.pool_arc()?;
25 Ok(Self { pool, event_sink })
26 }
27
28 pub async fn create_event(
29 &self,
30 session_id: &SessionId,
31 user_id: &UserId,
32 input: &CreateAnalyticsEventInput,
33 ) -> Result<AnalyticsEventCreated> {
34 let event = Self::build_record(session_id, user_id, input);
35 self.event_sink
36 .persist_events(std::slice::from_ref(&event))
37 .await?;
38 Ok(AnalyticsEventCreated {
39 id: event.id,
40 event_type: event.event_type,
41 })
42 }
43
44 pub async fn create_events_batch(
45 &self,
46 session_id: &SessionId,
47 user_id: &UserId,
48 inputs: &[CreateAnalyticsEventInput],
49 ) -> Result<Vec<AnalyticsEventCreated>> {
50 if inputs.is_empty() {
51 return Ok(Vec::new());
52 }
53
54 let events: Vec<_> = inputs
55 .iter()
56 .map(|input| Self::build_record(session_id, user_id, input))
57 .collect();
58 self.event_sink.persist_events(&events).await?;
59 Ok(events
60 .into_iter()
61 .map(|event| AnalyticsEventCreated {
62 id: event.id,
63 event_type: event.event_type,
64 })
65 .collect())
66 }
67
68 fn build_record(
69 session_id: &SessionId,
70 user_id: &UserId,
71 input: &CreateAnalyticsEventInput,
72 ) -> AnalyticsEventRecord {
73 AnalyticsEventRecord {
74 id: format!("evt_{}", uuid::Uuid::new_v4()),
75 user_id: user_id.clone(),
76 session_id: session_id.clone(),
77 event_type: input.event_type.as_str().to_owned(),
78 event_category: input.event_type.category().to_owned(),
79 page_url: input.page_url.clone(),
80 event_data: Self::build_event_data(input),
81 }
82 }
83
84 pub async fn count_events_by_type(
85 &self,
86 session_id: &SessionId,
87 event_type: &AnalyticsEventType,
88 ) -> Result<i64> {
89 let count = sqlx::query_scalar!(
90 r#"
91 SELECT COUNT(*) as "count!"
92 FROM analytics_report_analytics_events
93 WHERE session_id = $1 AND event_type = $2
94 "#,
95 session_id.as_str(),
96 event_type.as_str()
97 )
98 .fetch_one(&*self.pool)
99 .await?;
100
101 Ok(count)
102 }
103
104 pub async fn find_by_session(
105 &self,
106 session_id: &SessionId,
107 limit: i64,
108 ) -> Result<Vec<StoredAnalyticsEvent>> {
109 let events = sqlx::query_as!(
110 StoredAnalyticsEvent,
111 r#"
112 SELECT
113 id,
114 user_id as "user_id: UserId",
115 session_id as "session_id: SessionId",
116 event_type,
117 event_category,
118 endpoint as page_url,
119 event_data,
120 timestamp
121 FROM analytics_report_analytics_events
122 WHERE session_id = $1
123 ORDER BY timestamp DESC
124 LIMIT $2
125 "#,
126 session_id.as_str(),
127 limit
128 )
129 .fetch_all(&*self.pool)
130 .await?;
131
132 Ok(events)
133 }
134
135 pub async fn find_by_content(
136 &self,
137 content_id: &ContentId,
138 limit: i64,
139 ) -> Result<Vec<StoredAnalyticsEvent>> {
140 let events = sqlx::query_as!(
141 StoredAnalyticsEvent,
142 r#"
143 SELECT
144 id,
145 user_id as "user_id: UserId",
146 session_id as "session_id: SessionId",
147 event_type,
148 event_category,
149 endpoint as page_url,
150 event_data,
151 timestamp
152 FROM analytics_report_analytics_events
153 WHERE event_data->>'content_id' = $1
154 ORDER BY timestamp DESC
155 LIMIT $2
156 "#,
157 content_id.as_str(),
158 limit
159 )
160 .fetch_all(&*self.pool)
161 .await?;
162
163 Ok(events)
164 }
165
166 fn build_event_data(input: &CreateAnalyticsEventInput) -> serde_json::Value {
167 let mut data = input.data.clone().unwrap_or(serde_json::json!({}));
168
169 if let Some(obj) = data.as_object_mut() {
170 if let Some(content_id) = &input.content_id {
171 obj.insert(
172 "content_id".to_owned(),
173 serde_json::json!(content_id.as_str()),
174 );
175 }
176 if let Some(slug) = &input.slug {
177 obj.insert("slug".to_owned(), serde_json::json!(slug));
178 }
179 if let Some(referrer) = &input.referrer {
180 obj.insert("referrer".to_owned(), serde_json::json!(referrer));
181 }
182 }
183
184 data
185 }
186}
187
188#[derive(Debug, Clone, sqlx::FromRow)]
189pub struct StoredAnalyticsEvent {
190 pub id: String,
191 pub user_id: UserId,
192 pub session_id: Option<SessionId>,
193 pub event_type: String,
194 pub event_category: String,
195 pub page_url: Option<String>,
196 pub event_data: Option<serde_json::Value>,
197 pub timestamp: chrono::DateTime<chrono::Utc>,
198}