systemprompt-analytics 0.25.0

Analytics for systemprompt.io AI governance infrastructure. Session, agent, tool, and microdollar-precision cost attribution across the MCP governance pipeline.
Documentation
//! Traffic-source, geography, device, and bot analytics.
//!
//! [`TrafficAnalyticsRepository`] reads `user_sessions` to break sessions
//! down by referrer source, country, and device, and to classify human
//! versus bot traffic (including a user-agent-driven bot taxonomy). An
//! `engaged_only` flag restricts the human-facing breakdowns to sessions with
//! a landing page and at least one request.
//!
//! Copyright (c) systemprompt.io — Business Source License 1.1.
//! See <https://systemprompt.io> for licensing details.

use crate::Result;
use chrono::{DateTime, Utc};
use sqlx::PgPool;
use std::sync::Arc;
use systemprompt_database::DbPool;

use crate::models::cli::{
    BotTotalsRow, BotTypeRow, DeviceRow, GeoRow, TrafficNavigationRow, TrafficPageRow,
    TrafficSourceRow,
};

#[derive(Debug, Clone, Copy)]
pub struct PageQuery<'a> {
    pub start: DateTime<Utc>,
    pub end: DateTime<Utc>,
    pub limit: i64,
    pub engaged_only: bool,
    pub referrer: Option<&'a str>,
    pub path_prefix: Option<&'a str>,
}

#[derive(Debug, Clone, Copy)]
pub struct NavigationQuery<'a> {
    pub start: DateTime<Utc>,
    pub end: DateTime<Utc>,
    pub limit: i64,
    pub path_prefix: Option<&'a str>,
    pub internal_only: bool,
}

#[derive(Debug)]
pub struct TrafficAnalyticsRepository {
    pool: Arc<PgPool>,
}

impl TrafficAnalyticsRepository {
    pub fn new(db: &DbPool) -> Result<Self> {
        let pool = db.pool_arc()?;
        Ok(Self { pool })
    }

    pub async fn get_sources(
        &self,
        start: DateTime<Utc>,
        end: DateTime<Utc>,
        limit: i64,
        engaged_only: bool,
    ) -> Result<Vec<TrafficSourceRow>> {
        if engaged_only {
            sqlx::query_as!(
                TrafficSourceRow,
                r#"
                SELECT
                    COALESCE(referrer_source, 'direct') as "source",
                    COUNT(*)::bigint as "count!"
                FROM v_engaged_traffic
                WHERE started_at >= $1 AND started_at < $2
                GROUP BY referrer_source
                ORDER BY COUNT(*) DESC
                LIMIT $3
                "#,
                start,
                end,
                limit
            )
            .fetch_all(&*self.pool)
            .await
            .map_err(Into::into)
        } else {
            sqlx::query_as!(
                TrafficSourceRow,
                r#"
                SELECT
                    COALESCE(referrer_source, 'direct') as "source",
                    COUNT(*)::bigint as "count!"
                FROM v_clean_traffic
                WHERE started_at >= $1 AND started_at < $2
                GROUP BY referrer_source
                ORDER BY COUNT(*) DESC
                LIMIT $3
                "#,
                start,
                end,
                limit
            )
            .fetch_all(&*self.pool)
            .await
            .map_err(Into::into)
        }
    }

    pub async fn get_pages(&self, query: PageQuery<'_>) -> Result<Vec<TrafficPageRow>> {
        let PageQuery {
            start,
            end,
            limit,
            engaged_only,
            referrer,
            path_prefix,
        } = query;
        if engaged_only {
            sqlx::query_as!(
                TrafficPageRow,
                r#"
                SELECT
                    landing_page as "page",
                    COALESCE(referrer_source, 'direct') as "source",
                    COUNT(*)::bigint as "count!"
                FROM v_engaged_traffic
                WHERE started_at >= $1 AND started_at < $2
                  AND ($3::text IS NULL OR COALESCE(referrer_source, 'direct') = $3)
                  AND ($4::text IS NULL OR landing_page LIKE $4 || '%')
                GROUP BY landing_page, referrer_source
                ORDER BY COUNT(*) DESC
                LIMIT $5
                "#,
                start,
                end,
                referrer,
                path_prefix,
                limit
            )
            .fetch_all(&*self.pool)
            .await
            .map_err(Into::into)
        } else {
            sqlx::query_as!(
                TrafficPageRow,
                r#"
                SELECT
                    landing_page as "page",
                    COALESCE(referrer_source, 'direct') as "source",
                    COUNT(*)::bigint as "count!"
                FROM v_clean_traffic
                WHERE started_at >= $1 AND started_at < $2
                  AND landing_page IS NOT NULL
                  AND ($3::text IS NULL OR COALESCE(referrer_source, 'direct') = $3)
                  AND ($4::text IS NULL OR landing_page LIKE $4 || '%')
                GROUP BY landing_page, referrer_source
                ORDER BY COUNT(*) DESC
                LIMIT $5
                "#,
                start,
                end,
                referrer,
                path_prefix,
                limit
            )
            .fetch_all(&*self.pool)
            .await
            .map_err(Into::into)
        }
    }

    pub async fn get_navigation(
        &self,
        query: NavigationQuery<'_>,
    ) -> Result<Vec<TrafficNavigationRow>> {
        let NavigationQuery {
            start,
            end,
            limit,
            path_prefix,
            internal_only,
        } = query;
        if internal_only {
            sqlx::query_as!(
                TrafficNavigationRow,
                r#"
                SELECT
                    endpoint as "from_path",
                    event_data->>'target_url' as "to_path",
                    COUNT(*)::bigint as "count!"
                FROM analytics_events
                WHERE event_type = 'link_click'
                  AND timestamp >= $1 AND timestamp < $2
                  AND ($3::text IS NULL OR event_data->>'target_url' LIKE $3 || '%')
                  AND COALESCE(event_data->>'is_external', 'false') <> 'true'
                GROUP BY endpoint, event_data->>'target_url'
                ORDER BY COUNT(*) DESC
                LIMIT $4
                "#,
                start,
                end,
                path_prefix,
                limit
            )
            .fetch_all(&*self.pool)
            .await
            .map_err(Into::into)
        } else {
            sqlx::query_as!(
                TrafficNavigationRow,
                r#"
                SELECT
                    endpoint as "from_path",
                    event_data->>'target_url' as "to_path",
                    COUNT(*)::bigint as "count!"
                FROM analytics_events
                WHERE event_type = 'link_click'
                  AND timestamp >= $1 AND timestamp < $2
                  AND ($3::text IS NULL OR event_data->>'target_url' LIKE $3 || '%')
                GROUP BY endpoint, event_data->>'target_url'
                ORDER BY COUNT(*) DESC
                LIMIT $4
                "#,
                start,
                end,
                path_prefix,
                limit
            )
            .fetch_all(&*self.pool)
            .await
            .map_err(Into::into)
        }
    }

    pub async fn get_geo_breakdown(
        &self,
        start: DateTime<Utc>,
        end: DateTime<Utc>,
        limit: i64,
        engaged_only: bool,
    ) -> Result<Vec<GeoRow>> {
        if engaged_only {
            sqlx::query_as!(
                GeoRow,
                r#"
                SELECT
                    COALESCE(country, 'Unknown') as "country",
                    COUNT(*)::bigint as "count!"
                FROM v_engaged_traffic
                WHERE started_at >= $1 AND started_at < $2
                GROUP BY country
                ORDER BY COUNT(*) DESC
                LIMIT $3
                "#,
                start,
                end,
                limit
            )
            .fetch_all(&*self.pool)
            .await
            .map_err(Into::into)
        } else {
            sqlx::query_as!(
                GeoRow,
                r#"
                SELECT
                    COALESCE(country, 'Unknown') as "country",
                    COUNT(*)::bigint as "count!"
                FROM v_clean_traffic
                WHERE started_at >= $1 AND started_at < $2
                GROUP BY country
                ORDER BY COUNT(*) DESC
                LIMIT $3
                "#,
                start,
                end,
                limit
            )
            .fetch_all(&*self.pool)
            .await
            .map_err(Into::into)
        }
    }

    pub async fn get_device_breakdown(
        &self,
        start: DateTime<Utc>,
        end: DateTime<Utc>,
        limit: i64,
        engaged_only: bool,
    ) -> Result<Vec<DeviceRow>> {
        if engaged_only {
            sqlx::query_as!(
                DeviceRow,
                r#"
                SELECT
                    COALESCE(device_type, 'unknown') as "device",
                    COALESCE(browser, 'unknown') as "browser",
                    COUNT(*)::bigint as "count!"
                FROM v_engaged_traffic
                WHERE started_at >= $1 AND started_at < $2
                GROUP BY device_type, browser
                ORDER BY COUNT(*) DESC
                LIMIT $3
                "#,
                start,
                end,
                limit
            )
            .fetch_all(&*self.pool)
            .await
            .map_err(Into::into)
        } else {
            sqlx::query_as!(
                DeviceRow,
                r#"
                SELECT
                    COALESCE(device_type, 'unknown') as "device",
                    COALESCE(browser, 'unknown') as "browser",
                    COUNT(*)::bigint as "count!"
                FROM v_clean_traffic
                WHERE started_at >= $1 AND started_at < $2
                GROUP BY device_type, browser
                ORDER BY COUNT(*) DESC
                LIMIT $3
                "#,
                start,
                end,
                limit
            )
            .fetch_all(&*self.pool)
            .await
            .map_err(Into::into)
        }
    }

    pub async fn get_bot_totals(
        &self,
        start: DateTime<Utc>,
        end: DateTime<Utc>,
    ) -> Result<BotTotalsRow> {
        // Why: Partitions every session into exactly one bucket; the flag and
        // engagement predicates must mirror v_clean_traffic / v_engaged_traffic.
        sqlx::query_as!(
            BotTotalsRow,
            r#"
            SELECT
                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!",
                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!",
                COUNT(*) FILTER (WHERE is_bot = true OR is_ai_crawler = true OR is_scanner = true OR is_behavioral_bot = true)::bigint as "bot!"
            FROM user_sessions
            WHERE started_at >= $1 AND started_at < $2
            "#,
            start,
            end
        )
        .fetch_one(&*self.pool)
        .await
        .map_err(Into::into)
    }

    pub async fn get_bot_breakdown(
        &self,
        start: DateTime<Utc>,
        end: DateTime<Utc>,
    ) -> Result<Vec<BotTypeRow>> {
        sqlx::query_as!(
            BotTypeRow,
            r#"
            SELECT
                bot_type as "bot_type",
                COUNT(*)::bigint as "count!"
            FROM v_bot_sessions
            WHERE started_at >= $1 AND started_at < $2
            GROUP BY 1
            ORDER BY COUNT(*) DESC
            "#,
            start,
            end
        )
        .fetch_all(&*self.pool)
        .await
        .map_err(Into::into)
    }
}