systemprompt_content/repository/link/
analytics.rs1use 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}