Skip to main content

systemprompt_analytics/repository/
events.rs

1//! Raw analytics ingestion through the logging sink and event lookups.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use 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}