Skip to main content

systemprompt_analytics/repository/
content_analytics.rs

1//! Repository for per-content view/engagement aggregates.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use crate::Result;
7use chrono::{DateTime, Utc};
8use sqlx::PgPool;
9use std::sync::Arc;
10use systemprompt_database::DbPool;
11
12use systemprompt_identifiers::{ContentId, SourceId};
13
14use crate::models::reporting::{ContentStatsRow, ContentTrendRow, TopContentRow};
15
16#[derive(Debug)]
17pub struct ContentAnalyticsRepository {
18    pool: Arc<PgPool>,
19}
20
21impl ContentAnalyticsRepository {
22    pub fn new(db: &DbPool) -> Self {
23        let pool = db.pool();
24        Self { pool }
25    }
26
27    pub async fn get_top_content(
28        &self,
29        start: DateTime<Utc>,
30        end: DateTime<Utc>,
31        limit: i64,
32    ) -> Result<Vec<TopContentRow>> {
33        sqlx::query_as!(
34            TopContentRow,
35            r#"
36            WITH content_stats AS (
37                SELECT
38                    ee.content_id,
39                    COUNT(*)::bigint as total_views,
40                    COUNT(DISTINCT ee.session_id)::bigint as unique_visitors,
41                    (AVG(LEAST(ee.time_on_page_ms, 1800000)) / 1000.0)::float8 as avg_time_on_page_seconds
42                FROM engagement_events ee
43                INNER JOIN report_clean_traffic us ON ee.session_id = us.session_id
44                WHERE ee.created_at >= $1 AND ee.created_at < $2
45                    AND ee.content_id IS NOT NULL                GROUP BY ee.content_id
46            )
47            SELECT
48                cs.content_id as "content_id!: ContentId",
49                mc.slug as "slug?",
50                mc.title as "title?",
51                mc.source_id as "source_id?: SourceId",
52                cs.total_views as "total_views!",
53                cs.unique_visitors as "unique_visitors!",
54                cs.avg_time_on_page_seconds::float8 as "avg_time_on_page_seconds",
55                NULL::text as "trend_direction"
56            FROM content_stats cs
57            LEFT JOIN report_markdown_content mc ON cs.content_id = mc.id
58            ORDER BY cs.total_views DESC
59            LIMIT $3
60            "#,
61            start,
62            end,
63            limit
64        )
65        .fetch_all(&*self.pool)
66        .await
67        .map_err(Into::into)
68    }
69
70    pub async fn popular_content_ids(
71        &self,
72        source_id: &SourceId,
73        days: i32,
74        limit: i64,
75    ) -> Result<Vec<ContentId>> {
76        let rows: Vec<String> = sqlx::query_scalar!(
77            r#"
78            SELECT mc.id as "id!"
79            FROM report_markdown_content mc
80            LEFT JOIN report_analytics_events ae ON
81                ae.event_type = 'page_view'
82                AND ae.event_category = 'content'
83                AND ae.endpoint = 'GET /' || mc.source_id || '/' || mc.slug
84                AND ae.timestamp >= CURRENT_TIMESTAMP - ($2 || ' days')::INTERVAL
85            LEFT JOIN report_users u ON ae.user_id = u.id
86            WHERE mc.source_id = $1
87            GROUP BY mc.id, mc.published_at
88            ORDER BY COUNT(DISTINCT CASE
89                WHEN u.id IS NOT NULL AND u.is_bot = FALSE AND u.is_scanner = FALSE
90                THEN ae.user_id
91            END) DESC, mc.published_at DESC
92            LIMIT $3
93            "#,
94            source_id.as_str(),
95            days.to_string(),
96            limit
97        )
98        .fetch_all(&*self.pool)
99        .await?;
100
101        Ok(rows.into_iter().map(ContentId::new).collect())
102    }
103
104    pub async fn get_stats(
105        &self,
106        start: DateTime<Utc>,
107        end: DateTime<Utc>,
108    ) -> Result<ContentStatsRow> {
109        sqlx::query_as!(
110            ContentStatsRow,
111            r#"
112            SELECT
113                COUNT(*)::bigint as "total_views!",
114                COUNT(DISTINCT ee.session_id)::bigint as "unique_visitors!",
115                COALESCE(AVG(LEAST(ee.time_on_page_ms, 1800000)) / 1000.0, 0)::float8 as "avg_time_on_page_seconds",
116                COALESCE(AVG(ee.max_scroll_depth), 0)::float8 as "avg_scroll_depth",
117                COALESCE(SUM(ee.click_count), 0)::bigint as "total_clicks!"
118            FROM engagement_events ee
119            INNER JOIN report_clean_traffic us ON ee.session_id = us.session_id
120            WHERE ee.created_at >= $1 AND ee.created_at < $2            "#,
121            start,
122            end
123        )
124        .fetch_one(&*self.pool)
125        .await
126        .map_err(Into::into)
127    }
128
129    pub async fn get_content_for_trends(
130        &self,
131        start: DateTime<Utc>,
132        end: DateTime<Utc>,
133    ) -> Result<Vec<ContentTrendRow>> {
134        sqlx::query_as!(
135            ContentTrendRow,
136            r#"
137            WITH date_series AS (
138                SELECT generate_series(
139                    date_trunc('day', $1::timestamptz),
140                    date_trunc('day', $2::timestamptz) - interval '1 day',
141                    '1 day'::interval
142                ) as day
143            ),
144            daily_stats AS (
145                SELECT
146                    date_trunc('day', ee.created_at) as day,
147                    COUNT(*)::bigint as views,
148                    COUNT(DISTINCT ee.session_id)::bigint as unique_visitors
149                FROM engagement_events ee
150                INNER JOIN report_clean_traffic us ON ee.session_id = us.session_id
151                WHERE ee.created_at >= $1 AND ee.created_at < $2                GROUP BY date_trunc('day', ee.created_at)
152            )
153            SELECT
154                ds.day as "timestamp!",
155                COALESCE(s.views, 0)::bigint as "views!",
156                COALESCE(s.unique_visitors, 0)::bigint as "unique_visitors!"
157            FROM date_series ds
158            LEFT JOIN daily_stats s ON ds.day = s.day
159            ORDER BY ds.day
160            "#,
161            start,
162            end
163        )
164        .fetch_all(&*self.pool)
165        .await
166        .map_err(Into::into)
167    }
168}