systemprompt_analytics/repository/session_signals/
mod.rs1mod engagement_queries;
14mod geo;
15
16use std::sync::Arc;
17
18use chrono::{DateTime, Utc};
19use sqlx::PgPool;
20use systemprompt_database::DbPool;
21use systemprompt_identifiers::SessionId;
22use systemprompt_traits::{DynAnalyticsEventStore, DynContentCatalogStats, DynSessionStore};
23
24use crate::{AnalyticsError, Result};
25
26#[derive(Clone)]
27pub struct SessionSignalsRepository {
28 write_pool: Arc<PgPool>,
29 owner: DynSessionStore,
30 events: DynAnalyticsEventStore,
31 content: DynContentCatalogStats,
32}
33
34impl std::fmt::Debug for SessionSignalsRepository {
35 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
36 f.debug_struct("SessionSignalsRepository")
37 .finish_non_exhaustive()
38 }
39}
40
41impl SessionSignalsRepository {
42 pub fn new(
43 db: &DbPool,
44 owner: DynSessionStore,
45 events: DynAnalyticsEventStore,
46 content: DynContentCatalogStats,
47 ) -> Self {
48 Self {
49 write_pool: db.write_pool(),
50 owner,
51 events,
52 content,
53 }
54 }
55
56 pub async fn get_endpoint_sequence(&self, session_id: &SessionId) -> Result<Vec<String>> {
57 self.events
58 .get_endpoint_sequence(session_id)
59 .await
60 .map_err(AnalyticsError::from)
61 }
62
63 pub async fn get_request_timestamps(
64 &self,
65 session_id: &SessionId,
66 ) -> Result<Vec<DateTime<Utc>>> {
67 self.events
68 .get_request_timestamps(session_id)
69 .await
70 .map_err(AnalyticsError::from)
71 }
72
73 pub async fn has_analytics_events(&self, session_id: &SessionId) -> Result<bool> {
74 self.events
75 .has_analytics_events(session_id)
76 .await
77 .map_err(AnalyticsError::from)
78 }
79
80 pub async fn get_total_content_pages(&self) -> Result<i64> {
81 self.content.count_public_pages().await.map_err(Into::into)
82 }
83
84 pub async fn count_engagement_events_by_fingerprint(
85 &self,
86 fingerprint_hash: &str,
87 window_days: i64,
88 ) -> Result<i64> {
89 let session_ids = self
90 .owner
91 .fingerprint_session_ids(fingerprint_hash, window_days)
92 .await
93 .map_err(AnalyticsError::from)?
94 .into_iter()
95 .map(|id| id.to_string())
96 .collect::<Vec<_>>();
97 engagement_queries::count_engagement_events_for_sessions(&self.write_pool, &session_ids)
98 .await
99 }
100
101 pub async fn count_sessions_missing_geo(&self) -> Result<i64> {
102 self.owner
103 .count_sessions_missing_geo()
104 .await
105 .map_err(AnalyticsError::from)
106 }
107}