Skip to main content

systemprompt_analytics/repository/
traffic.rs

1//! Traffic-source, geography, device, and bot analytics.
2//!
3//! [`TrafficAnalyticsRepository`] reads `user_sessions` to break sessions
4//! down by referrer source, country, and device, and to classify human
5//! versus bot traffic (including a user-agent-driven bot taxonomy). An
6//! `engaged_only` flag restricts the human-facing breakdowns to sessions with
7//! a landing page and at least one request.
8//!
9//! Copyright (c) systemprompt.io — Business Source License 1.1.
10//! See <https://systemprompt.io> for licensing details.
11
12use crate::Result;
13use chrono::{DateTime, Utc};
14use sqlx::PgPool;
15use std::sync::Arc;
16use systemprompt_database::DbPool;
17
18use crate::models::cli::{
19    BotTotalsRow, BotTypeRow, DeviceRow, GeoRow, TrafficNavigationRow, TrafficPageRow,
20    TrafficSourceRow,
21};
22
23#[derive(Debug, Clone, Copy)]
24pub struct PageQuery<'a> {
25    pub start: DateTime<Utc>,
26    pub end: DateTime<Utc>,
27    pub limit: i64,
28    pub engaged_only: bool,
29    pub referrer: Option<&'a str>,
30    pub path_prefix: Option<&'a str>,
31}
32
33#[derive(Debug, Clone, Copy)]
34pub struct NavigationQuery<'a> {
35    pub start: DateTime<Utc>,
36    pub end: DateTime<Utc>,
37    pub limit: i64,
38    pub path_prefix: Option<&'a str>,
39    pub internal_only: bool,
40}
41
42#[derive(Debug)]
43pub struct TrafficAnalyticsRepository {
44    pool: Arc<PgPool>,
45}
46
47impl TrafficAnalyticsRepository {
48    pub fn new(db: &DbPool) -> Result<Self> {
49        let pool = db.pool_arc()?;
50        Ok(Self { pool })
51    }
52
53    pub async fn get_sources(
54        &self,
55        start: DateTime<Utc>,
56        end: DateTime<Utc>,
57        limit: i64,
58        engaged_only: bool,
59    ) -> Result<Vec<TrafficSourceRow>> {
60        if engaged_only {
61            sqlx::query_as!(
62                TrafficSourceRow,
63                r#"
64                SELECT
65                    COALESCE(referrer_source, 'direct') as "source",
66                    COUNT(*)::bigint as "count!"
67                FROM v_engaged_traffic
68                WHERE started_at >= $1 AND started_at < $2
69                GROUP BY referrer_source
70                ORDER BY COUNT(*) DESC
71                LIMIT $3
72                "#,
73                start,
74                end,
75                limit
76            )
77            .fetch_all(&*self.pool)
78            .await
79            .map_err(Into::into)
80        } else {
81            sqlx::query_as!(
82                TrafficSourceRow,
83                r#"
84                SELECT
85                    COALESCE(referrer_source, 'direct') as "source",
86                    COUNT(*)::bigint as "count!"
87                FROM v_clean_traffic
88                WHERE started_at >= $1 AND started_at < $2
89                GROUP BY referrer_source
90                ORDER BY COUNT(*) DESC
91                LIMIT $3
92                "#,
93                start,
94                end,
95                limit
96            )
97            .fetch_all(&*self.pool)
98            .await
99            .map_err(Into::into)
100        }
101    }
102
103    pub async fn get_pages(&self, query: PageQuery<'_>) -> Result<Vec<TrafficPageRow>> {
104        let PageQuery {
105            start,
106            end,
107            limit,
108            engaged_only,
109            referrer,
110            path_prefix,
111        } = query;
112        if engaged_only {
113            sqlx::query_as!(
114                TrafficPageRow,
115                r#"
116                SELECT
117                    landing_page as "page",
118                    COALESCE(referrer_source, 'direct') as "source",
119                    COUNT(*)::bigint as "count!"
120                FROM v_engaged_traffic
121                WHERE started_at >= $1 AND started_at < $2
122                  AND ($3::text IS NULL OR COALESCE(referrer_source, 'direct') = $3)
123                  AND ($4::text IS NULL OR landing_page LIKE $4 || '%')
124                GROUP BY landing_page, referrer_source
125                ORDER BY COUNT(*) DESC
126                LIMIT $5
127                "#,
128                start,
129                end,
130                referrer,
131                path_prefix,
132                limit
133            )
134            .fetch_all(&*self.pool)
135            .await
136            .map_err(Into::into)
137        } else {
138            sqlx::query_as!(
139                TrafficPageRow,
140                r#"
141                SELECT
142                    landing_page as "page",
143                    COALESCE(referrer_source, 'direct') as "source",
144                    COUNT(*)::bigint as "count!"
145                FROM v_clean_traffic
146                WHERE started_at >= $1 AND started_at < $2
147                  AND landing_page IS NOT NULL
148                  AND ($3::text IS NULL OR COALESCE(referrer_source, 'direct') = $3)
149                  AND ($4::text IS NULL OR landing_page LIKE $4 || '%')
150                GROUP BY landing_page, referrer_source
151                ORDER BY COUNT(*) DESC
152                LIMIT $5
153                "#,
154                start,
155                end,
156                referrer,
157                path_prefix,
158                limit
159            )
160            .fetch_all(&*self.pool)
161            .await
162            .map_err(Into::into)
163        }
164    }
165
166    pub async fn get_navigation(
167        &self,
168        query: NavigationQuery<'_>,
169    ) -> Result<Vec<TrafficNavigationRow>> {
170        let NavigationQuery {
171            start,
172            end,
173            limit,
174            path_prefix,
175            internal_only,
176        } = query;
177        if internal_only {
178            sqlx::query_as!(
179                TrafficNavigationRow,
180                r#"
181                SELECT
182                    endpoint as "from_path",
183                    event_data->>'target_url' as "to_path",
184                    COUNT(*)::bigint as "count!"
185                FROM analytics_events
186                WHERE event_type = 'link_click'
187                  AND timestamp >= $1 AND timestamp < $2
188                  AND ($3::text IS NULL OR event_data->>'target_url' LIKE $3 || '%')
189                  AND COALESCE(event_data->>'is_external', 'false') <> 'true'
190                GROUP BY endpoint, event_data->>'target_url'
191                ORDER BY COUNT(*) DESC
192                LIMIT $4
193                "#,
194                start,
195                end,
196                path_prefix,
197                limit
198            )
199            .fetch_all(&*self.pool)
200            .await
201            .map_err(Into::into)
202        } else {
203            sqlx::query_as!(
204                TrafficNavigationRow,
205                r#"
206                SELECT
207                    endpoint as "from_path",
208                    event_data->>'target_url' as "to_path",
209                    COUNT(*)::bigint as "count!"
210                FROM analytics_events
211                WHERE event_type = 'link_click'
212                  AND timestamp >= $1 AND timestamp < $2
213                  AND ($3::text IS NULL OR event_data->>'target_url' LIKE $3 || '%')
214                GROUP BY endpoint, event_data->>'target_url'
215                ORDER BY COUNT(*) DESC
216                LIMIT $4
217                "#,
218                start,
219                end,
220                path_prefix,
221                limit
222            )
223            .fetch_all(&*self.pool)
224            .await
225            .map_err(Into::into)
226        }
227    }
228
229    pub async fn get_geo_breakdown(
230        &self,
231        start: DateTime<Utc>,
232        end: DateTime<Utc>,
233        limit: i64,
234        engaged_only: bool,
235    ) -> Result<Vec<GeoRow>> {
236        if engaged_only {
237            sqlx::query_as!(
238                GeoRow,
239                r#"
240                SELECT
241                    COALESCE(country, 'Unknown') as "country",
242                    COUNT(*)::bigint as "count!"
243                FROM v_engaged_traffic
244                WHERE started_at >= $1 AND started_at < $2
245                GROUP BY country
246                ORDER BY COUNT(*) DESC
247                LIMIT $3
248                "#,
249                start,
250                end,
251                limit
252            )
253            .fetch_all(&*self.pool)
254            .await
255            .map_err(Into::into)
256        } else {
257            sqlx::query_as!(
258                GeoRow,
259                r#"
260                SELECT
261                    COALESCE(country, 'Unknown') as "country",
262                    COUNT(*)::bigint as "count!"
263                FROM v_clean_traffic
264                WHERE started_at >= $1 AND started_at < $2
265                GROUP BY country
266                ORDER BY COUNT(*) DESC
267                LIMIT $3
268                "#,
269                start,
270                end,
271                limit
272            )
273            .fetch_all(&*self.pool)
274            .await
275            .map_err(Into::into)
276        }
277    }
278
279    pub async fn get_device_breakdown(
280        &self,
281        start: DateTime<Utc>,
282        end: DateTime<Utc>,
283        limit: i64,
284        engaged_only: bool,
285    ) -> Result<Vec<DeviceRow>> {
286        if engaged_only {
287            sqlx::query_as!(
288                DeviceRow,
289                r#"
290                SELECT
291                    COALESCE(device_type, 'unknown') as "device",
292                    COALESCE(browser, 'unknown') as "browser",
293                    COUNT(*)::bigint as "count!"
294                FROM v_engaged_traffic
295                WHERE started_at >= $1 AND started_at < $2
296                GROUP BY device_type, browser
297                ORDER BY COUNT(*) DESC
298                LIMIT $3
299                "#,
300                start,
301                end,
302                limit
303            )
304            .fetch_all(&*self.pool)
305            .await
306            .map_err(Into::into)
307        } else {
308            sqlx::query_as!(
309                DeviceRow,
310                r#"
311                SELECT
312                    COALESCE(device_type, 'unknown') as "device",
313                    COALESCE(browser, 'unknown') as "browser",
314                    COUNT(*)::bigint as "count!"
315                FROM v_clean_traffic
316                WHERE started_at >= $1 AND started_at < $2
317                GROUP BY device_type, browser
318                ORDER BY COUNT(*) DESC
319                LIMIT $3
320                "#,
321                start,
322                end,
323                limit
324            )
325            .fetch_all(&*self.pool)
326            .await
327            .map_err(Into::into)
328        }
329    }
330
331    pub async fn get_bot_totals(
332        &self,
333        start: DateTime<Utc>,
334        end: DateTime<Utc>,
335    ) -> Result<BotTotalsRow> {
336        // Why: Partitions every session into exactly one bucket; the flag and
337        // engagement predicates must mirror v_clean_traffic / v_engaged_traffic.
338        sqlx::query_as!(
339            BotTotalsRow,
340            r#"
341            SELECT
342                COUNT(*) FILTER (WHERE is_bot = false AND is_ai_crawler = false AND is_scanner = false AND is_behavioral_bot = false AND landing_page IS NOT NULL AND request_count > 0)::bigint as "human!",
343                COUNT(*) FILTER (WHERE is_bot = false AND is_ai_crawler = false AND is_scanner = false AND is_behavioral_bot = false AND (landing_page IS NULL OR request_count = 0))::bigint as "ghost!",
344                COUNT(*) FILTER (WHERE is_bot = true OR is_ai_crawler = true OR is_scanner = true OR is_behavioral_bot = true)::bigint as "bot!"
345            FROM user_sessions
346            WHERE started_at >= $1 AND started_at < $2
347            "#,
348            start,
349            end
350        )
351        .fetch_one(&*self.pool)
352        .await
353        .map_err(Into::into)
354    }
355
356    pub async fn get_bot_breakdown(
357        &self,
358        start: DateTime<Utc>,
359        end: DateTime<Utc>,
360    ) -> Result<Vec<BotTypeRow>> {
361        sqlx::query_as!(
362            BotTypeRow,
363            r#"
364            SELECT
365                bot_type as "bot_type",
366                COUNT(*)::bigint as "count!"
367            FROM v_bot_sessions
368            WHERE started_at >= $1 AND started_at < $2
369            GROUP BY 1
370            ORDER BY COUNT(*) DESC
371            "#,
372            start,
373            end
374        )
375        .fetch_all(&*self.pool)
376        .await
377        .map_err(Into::into)
378    }
379}