Skip to main content

systemprompt_content/repository/link/
analytics.rs

1//! Link click and analytics repository.
2//!
3//! [`LinkAnalyticsRepository`] records click events and serves the aggregate
4//! click/conversion views over `link_clicks` and `campaign_links`, maintaining
5//! the denormalised counters on the link row as clicks arrive.
6
7use crate::error::ContentError;
8use crate::models::{
9    CampaignPerformance, ContentJourneyNode, LinkClick, LinkPerformance, RecordClickParams,
10};
11use sqlx::PgPool;
12use std::sync::Arc;
13use systemprompt_database::DbPool;
14use systemprompt_identifiers::{
15    CampaignId, ContentId, ContextId, LinkClickId, LinkId, SessionId, TaskId, UserId,
16};
17
18#[derive(Debug)]
19pub struct LinkAnalyticsRepository {
20    pool: Arc<PgPool>,
21    write_pool: Arc<PgPool>,
22}
23
24impl LinkAnalyticsRepository {
25    pub fn new(db: &DbPool) -> Result<Self, ContentError> {
26        let pool = db
27            .pool_arc()
28            .map_err(|e| ContentError::InvalidRequest(format!("Database pool error: {e}")))?;
29        let write_pool = db
30            .write_pool_arc()
31            .map_err(|e| ContentError::InvalidRequest(format!("Database write pool error: {e}")))?;
32        Ok(Self { pool, write_pool })
33    }
34
35    pub async fn get_link_performance(
36        &self,
37        link_id: &LinkId,
38    ) -> Result<Option<LinkPerformance>, sqlx::Error> {
39        sqlx::query_as!(
40            LinkPerformance,
41            r#"
42            SELECT
43                l.id as "link_id: LinkId",
44                COALESCE(l.click_count, 0)::bigint as "click_count!",
45                COALESCE(l.unique_click_count, 0)::bigint as "unique_click_count!",
46                COALESCE(l.conversion_count, 0)::bigint as "conversion_count!",
47                CASE
48                    WHEN COALESCE(l.click_count, 0) > 0 THEN
49                        COALESCE(l.conversion_count, 0)::float / l.click_count
50                    ELSE 0.0
51                END as conversion_rate
52            FROM campaign_links l
53            WHERE l.id = $1
54            "#,
55            link_id.as_str()
56        )
57        .fetch_optional(&*self.pool)
58        .await
59    }
60
61    pub async fn check_session_clicked_link(
62        &self,
63        link_id: &LinkId,
64        session_id: &SessionId,
65    ) -> Result<bool, sqlx::Error> {
66        let result = sqlx::query!(
67            r#"SELECT COALESCE(COUNT(*), 0)::bigint as "count!" FROM link_clicks WHERE link_id = $1 AND session_id = $2"#,
68            link_id.as_str(),
69            session_id.as_str()
70        )
71        .fetch_one(&*self.pool)
72        .await?;
73
74        Ok(result.count > 0)
75    }
76
77    pub async fn increment_link_clicks(
78        &self,
79        link_id: &LinkId,
80        is_first_click: bool,
81    ) -> Result<(), sqlx::Error> {
82        if is_first_click {
83            sqlx::query!(
84                "UPDATE campaign_links SET click_count = click_count + 1, unique_click_count = \
85                 unique_click_count + 1 WHERE id = $1",
86                link_id.as_str()
87            )
88            .execute(&*self.write_pool)
89            .await?;
90        } else {
91            sqlx::query!(
92                "UPDATE campaign_links SET click_count = click_count + 1 WHERE id = $1",
93                link_id.as_str()
94            )
95            .execute(&*self.write_pool)
96            .await?;
97        }
98        Ok(())
99    }
100
101    pub async fn get_clicks_by_link(
102        &self,
103        link_id: &LinkId,
104        limit: i64,
105        offset: i64,
106    ) -> Result<Vec<LinkClick>, sqlx::Error> {
107        sqlx::query_as!(
108            LinkClick,
109            r#"
110            SELECT id as "id: LinkClickId", link_id as "link_id: LinkId",
111                   session_id as "session_id: SessionId", user_id as "user_id: UserId",
112                   context_id as "context_id: ContextId", task_id as "task_id: TaskId",
113                   referrer_page, referrer_url, clicked_at, user_agent, ip_address,
114                   device_type, country, is_first_click, is_conversion, conversion_at,
115                   time_on_page_seconds, scroll_depth_percent
116            FROM link_clicks
117            WHERE link_id = $1
118            ORDER BY clicked_at DESC
119            LIMIT $2 OFFSET $3
120            "#,
121            link_id.as_str(),
122            limit,
123            offset
124        )
125        .fetch_all(&*self.pool)
126        .await
127    }
128
129    pub async fn get_content_journey_map(
130        &self,
131        limit: i64,
132        offset: i64,
133    ) -> Result<Vec<ContentJourneyNode>, sqlx::Error> {
134        let rows = sqlx::query!(
135            r#"
136            SELECT source_content_id, target_url, COALESCE(click_count, 0) as "click_count!"
137            FROM campaign_links
138            WHERE source_content_id IS NOT NULL AND click_count > 0
139            ORDER BY click_count DESC
140            LIMIT $1 OFFSET $2
141            "#,
142            limit,
143            offset
144        )
145        .fetch_all(&*self.pool)
146        .await?;
147
148        Ok(rows
149            .into_iter()
150            .filter_map(|r| {
151                Some(ContentJourneyNode {
152                    source_content_id: ContentId::new(r.source_content_id?),
153                    target_url: r.target_url,
154                    click_count: r.click_count,
155                })
156            })
157            .collect())
158    }
159
160    pub async fn get_campaign_performance(
161        &self,
162        campaign_id: &CampaignId,
163    ) -> Result<Option<CampaignPerformance>, sqlx::Error> {
164        sqlx::query_as!(
165            CampaignPerformance,
166            r#"
167            SELECT
168                campaign_id as "campaign_id!: CampaignId",
169                COALESCE(SUM(click_count), 0)::bigint as "total_clicks!",
170                COUNT(*)::bigint as "link_count!",
171                COUNT(DISTINCT source_content_id) as unique_visitors,
172                COALESCE(SUM(conversion_count), 0)::bigint as conversion_count
173            FROM campaign_links
174            WHERE campaign_id = $1
175            GROUP BY campaign_id
176            "#,
177            campaign_id.as_str()
178        )
179        .fetch_optional(&*self.pool)
180        .await
181    }
182
183    pub async fn record_click(&self, params: &RecordClickParams) -> Result<(), sqlx::Error> {
184        sqlx::query!(
185            r#"
186            INSERT INTO link_clicks (
187                id, link_id, session_id, user_id, context_id, task_id,
188                referrer_page, referrer_url, clicked_at, user_agent, ip_address,
189                device_type, country, is_first_click, is_conversion
190            )
191            VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15)
192            "#,
193            params.click_id.as_str(),
194            params.link_id.as_str(),
195            params.session_id.as_str(),
196            params.user_id.as_ref().map(UserId::as_str),
197            params.context_id.as_ref().map(ContextId::as_str),
198            params.task_id.as_ref().map(TaskId::as_str),
199            params.referrer_page,
200            params.referrer_url,
201            params.clicked_at,
202            params.user_agent,
203            params.ip_address,
204            params.device_type,
205            params.country,
206            params.is_first_click,
207            params.is_conversion
208        )
209        .execute(&*self.write_pool)
210        .await?;
211        Ok(())
212    }
213}