Skip to main content

systemprompt_analytics/repository/session_signals/
mod.rs

1//! Session signals analytics composes across owners.
2//!
3//! Plain session reads and writes belong to the users domain and are reached
4//! through [`systemprompt_traits::SessionStore`] directly. What lives here is
5//! only what analytics adds on top: behavioural-detector inputs drawn from the
6//! analytics event store and the content catalogue, engagement counts over the
7//! sessions a fingerprint owns, and geolocation backfill for sessions that
8//! arrived without it.
9//!
10//! Copyright (c) systemprompt.io — Business Source License 1.1.
11//! See <https://systemprompt.io> for licensing details.
12
13mod 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}