Skip to main content

systemprompt_analytics/repository/
engagement.rs

1//! Persistence for page-level engagement telemetry.
2//!
3//! [`EngagementRepository`] records [`EngagementEvent`]s (scroll, click,
4//! focus, and reading-pattern metrics) and reads them back by id or per
5//! user. Writes go to the write pool; reads to the read pool.
6//!
7//! Copyright (c) systemprompt.io — Business Source License 1.1.
8//! See <https://systemprompt.io> for licensing details.
9
10use std::sync::Arc;
11
12use crate::Result;
13use sqlx::PgPool;
14use systemprompt_database::DbPool;
15use systemprompt_identifiers::{ContentId, EngagementEventId, SessionId, UserId};
16
17use crate::models::{CreateEngagementEventInput, EngagementEvent, EngagementEventRow};
18
19#[derive(Clone, Debug)]
20pub struct EngagementRepository {
21    pool: Arc<PgPool>,
22    write_pool: Arc<PgPool>,
23}
24
25impl EngagementRepository {
26    pub fn new(db: &DbPool) -> Self {
27        let pool = db.pool();
28        let write_pool = db.write_pool();
29        Self { pool, write_pool }
30    }
31
32    pub async fn create_engagement(
33        &self,
34        session_id: &SessionId,
35        user_id: &UserId,
36        content_id: Option<&ContentId>,
37        input: &CreateEngagementEventInput,
38    ) -> Result<EngagementEventId> {
39        let id = EngagementEventId::generate();
40
41        sqlx::query!(
42            r#"
43            INSERT INTO engagement_events (
44                id, session_id, user_id, page_url, content_id, event_type,
45                time_on_page_ms, max_scroll_depth, click_count,
46                time_to_first_interaction_ms, time_to_first_scroll_ms,
47                scroll_velocity_avg, scroll_direction_changes,
48                mouse_move_distance_px, keyboard_events, copy_events,
49                focus_time_ms, blur_count, tab_switches, visible_time_ms, hidden_time_ms,
50                is_rage_click, is_dead_click, reading_pattern, event_data
51            )
52            VALUES (
53                $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13,
54                $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $25
55            )
56            "#,
57            id.as_str(),
58            session_id.as_str(),
59            user_id.as_str(),
60            input.page_url,
61            content_id.map(ContentId::as_str),
62            input.event_type.as_str(),
63            input.time_on_page_ms,
64            input.max_scroll_depth,
65            input.click_count,
66            input.optional_metrics.time_to_first_interaction_ms,
67            input.optional_metrics.time_to_first_scroll_ms,
68            input.optional_metrics.scroll_velocity_avg,
69            input.optional_metrics.scroll_direction_changes,
70            input.optional_metrics.mouse_move_distance_px,
71            input.optional_metrics.keyboard_events,
72            input.optional_metrics.copy_events,
73            input.optional_metrics.focus_time_ms.unwrap_or(0),
74            input.optional_metrics.blur_count.unwrap_or(0),
75            input.optional_metrics.tab_switches.unwrap_or(0),
76            input.optional_metrics.visible_time_ms.unwrap_or(0),
77            input.optional_metrics.hidden_time_ms.unwrap_or(0),
78            input.optional_metrics.is_rage_click,
79            input.optional_metrics.is_dead_click,
80            input.optional_metrics.reading_pattern,
81            input.event_data.clone()
82        )
83        .execute(&*self.write_pool)
84        .await?;
85
86        Ok(id)
87    }
88
89    pub async fn find_by_id(&self, id: &EngagementEventId) -> Result<Option<EngagementEvent>> {
90        let event = sqlx::query_as!(
91            EngagementEventRow,
92            r#"
93            SELECT
94                id as "id: EngagementEventId", session_id, user_id, page_url,
95                content_id as "content_id: ContentId",
96                event_type,
97                time_on_page_ms, time_to_first_interaction_ms, time_to_first_scroll_ms,
98                max_scroll_depth, scroll_velocity_avg, scroll_direction_changes,
99                click_count, mouse_move_distance_px, keyboard_events, copy_events,
100                focus_time_ms as "focus_time_ms!",
101                blur_count as "blur_count!",
102                tab_switches as "tab_switches!",
103                visible_time_ms as "visible_time_ms!",
104                hidden_time_ms as "hidden_time_ms!",
105                is_rage_click, is_dead_click, reading_pattern,
106                created_at, updated_at
107            FROM engagement_events
108            WHERE id = $1
109            "#,
110            id.as_str()
111        )
112        .fetch_optional(&*self.pool)
113        .await
114        .map(|row| row.map(EngagementEvent::from))?;
115
116        Ok(event)
117    }
118
119    pub async fn list_by_user(&self, user_id: &UserId, limit: i64) -> Result<Vec<EngagementEvent>> {
120        let events = sqlx::query_as!(
121            EngagementEventRow,
122            r#"
123            SELECT
124                id as "id: EngagementEventId", session_id, user_id, page_url,
125                content_id as "content_id: ContentId",
126                event_type,
127                time_on_page_ms as "time_on_page_ms!", time_to_first_interaction_ms, time_to_first_scroll_ms,
128                max_scroll_depth as "max_scroll_depth!", scroll_velocity_avg, scroll_direction_changes,
129                click_count as "click_count!", mouse_move_distance_px, keyboard_events, copy_events,
130                focus_time_ms as "focus_time_ms!",
131                blur_count as "blur_count!",
132                tab_switches as "tab_switches!",
133                visible_time_ms as "visible_time_ms!",
134                hidden_time_ms as "hidden_time_ms!",
135                is_rage_click, is_dead_click, reading_pattern,
136                created_at, updated_at
137            FROM engagement_events
138            WHERE user_id = $1
139            ORDER BY created_at DESC
140            LIMIT $2
141            "#,
142            user_id.as_str(),
143            limit
144        )
145        .fetch_all(&*self.pool)
146        .await.map(|rows| rows.into_iter().map(EngagementEvent::from).collect::<Vec<_>>())?;
147
148        Ok(events)
149    }
150}