Skip to main content

systemprompt_analytics/repository/
cli_sessions.rs

1//! Session analytics over the `user_sessions` table.
2//!
3//! [`CliSessionAnalyticsRepository`] reads session counts, durations, live
4//! activity, and conversion stats for human (non-bot) sessions; every query
5//! filters out bot, behavioural-bot, and scanner traffic.
6//!
7//! Copyright (c) systemprompt.io — Business Source License 1.1.
8//! See <https://systemprompt.io> for licensing details.
9
10use crate::Result;
11use chrono::{DateTime, Utc};
12use sqlx::PgPool;
13use std::sync::Arc;
14use systemprompt_database::DbPool;
15use systemprompt_identifiers::UserId;
16
17use crate::models::cli::{LiveSessionRow, SessionStatsRow, SessionTrendRow};
18
19#[derive(Debug)]
20pub struct CliSessionAnalyticsRepository {
21    pool: Arc<PgPool>,
22}
23
24impl CliSessionAnalyticsRepository {
25    pub fn new(db: &DbPool) -> Result<Self> {
26        let pool = db.pool_arc()?;
27        Ok(Self { pool })
28    }
29
30    pub async fn get_stats(
31        &self,
32        start: DateTime<Utc>,
33        end: DateTime<Utc>,
34    ) -> Result<SessionStatsRow> {
35        sqlx::query_as!(
36            SessionStatsRow,
37            r#"
38            SELECT
39                COUNT(*)::bigint as "total_sessions!",
40                COUNT(DISTINCT user_id)::bigint as "unique_users!",
41                AVG(LEAST(
42                    COALESCE(
43                        duration_seconds,
44                        EXTRACT(EPOCH FROM (COALESCE(ended_at, last_activity_at) - started_at))::INTEGER
45                    ),
46                    1800
47                ))::float8 as "avg_duration",
48                AVG(request_count)::float8 as "avg_requests",
49                COUNT(*) FILTER (WHERE converted_at IS NOT NULL)::bigint as "conversions!"
50            FROM v_clean_traffic
51            WHERE started_at >= $1 AND started_at < $2            "#,
52            start,
53            end
54        )
55        .fetch_one(&*self.pool)
56        .await
57        .map_err(Into::into)
58    }
59
60    pub async fn get_active_session_count(&self, since: DateTime<Utc>) -> Result<i64> {
61        let count = sqlx::query_scalar!(
62            r#"SELECT COUNT(*)::bigint as "count!" FROM v_clean_traffic WHERE ended_at IS NULL AND last_activity_at >= $1"#,
63            since
64        )
65        .fetch_one(&*self.pool)
66        .await?;
67        Ok(count)
68    }
69
70    pub async fn get_live_sessions(
71        &self,
72        cutoff: DateTime<Utc>,
73        limit: i64,
74    ) -> Result<Vec<LiveSessionRow>> {
75        sqlx::query_as!(
76            LiveSessionRow,
77            r#"
78            SELECT
79                session_id as "session_id!",
80                COALESCE(user_type, 'unknown') as "user_type",
81                started_at as "started_at!",
82                duration_seconds,
83                request_count,
84                last_activity_at as "last_activity_at!"
85            FROM v_clean_traffic
86            WHERE ended_at IS NULL
87              AND last_activity_at >= $1            ORDER BY last_activity_at DESC
88            LIMIT $2
89            "#,
90            cutoff,
91            limit
92        )
93        .fetch_all(&*self.pool)
94        .await
95        .map_err(Into::into)
96    }
97
98    pub async fn get_active_count(&self, cutoff: DateTime<Utc>) -> Result<i64> {
99        let count = sqlx::query_scalar!(
100            r#"SELECT COUNT(*)::bigint as "count!" FROM v_clean_traffic WHERE ended_at IS NULL AND last_activity_at >= $1"#,
101            cutoff
102        )
103        .fetch_one(&*self.pool)
104        .await?;
105        Ok(count)
106    }
107
108    pub async fn get_sessions_for_trends(
109        &self,
110        start: DateTime<Utc>,
111        end: DateTime<Utc>,
112    ) -> Result<Vec<SessionTrendRow>> {
113        sqlx::query_as!(
114            SessionTrendRow,
115            r#"
116            SELECT
117                started_at as "started_at!",
118                user_id as "user_id: UserId",
119                duration_seconds
120            FROM v_clean_traffic
121            WHERE started_at >= $1 AND started_at < $2            ORDER BY started_at
122            "#,
123            start,
124            end
125        )
126        .fetch_all(&*self.pool)
127        .await
128        .map_err(Into::into)
129    }
130
131    pub async fn get_active_count_since(&self, start: DateTime<Utc>) -> Result<i64> {
132        let count = sqlx::query_scalar!(
133            r#"
134            SELECT COUNT(*)::bigint as "count!"
135            FROM v_clean_traffic
136            WHERE ended_at IS NULL
137              AND last_activity_at >= $1            "#,
138            start
139        )
140        .fetch_one(&*self.pool)
141        .await?;
142        Ok(count)
143    }
144
145    pub async fn get_total_count(&self, start: DateTime<Utc>, end: DateTime<Utc>) -> Result<i64> {
146        let count = sqlx::query_scalar!(
147            r#"SELECT COUNT(*)::bigint as "count!" FROM v_clean_traffic WHERE started_at >= $1 AND started_at < $2"#,
148            start,
149            end
150        )
151        .fetch_one(&*self.pool)
152        .await?;
153        Ok(count)
154    }
155}