Skip to main content

systemprompt_users/sessions/
mod.rs

1//! Authoritative session persistence contracts.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6mod behavioral;
7mod behavioral_queries;
8mod geo;
9mod mutations;
10mod providers;
11mod queries;
12use crate::Result;
13use chrono::{DateTime, Utc};
14use sqlx::PgPool;
15use std::sync::Arc;
16use systemprompt_database::DbPool;
17use systemprompt_identifiers::{SessionId, UserId};
18use systemprompt_traits::session_store::{
19    ActiveSessionLookup, CreateSessionParams, SessionBehavioralData, SessionRecord,
20    SessionSnapshot as AnalyticsSession,
21};
22#[derive(Clone, Debug)]
23pub struct SessionRepository {
24    pool: Arc<PgPool>,
25    write_pool: Arc<PgPool>,
26}
27impl SessionRepository {
28    pub fn new(db: &DbPool) -> Result<Self> {
29        Ok(Self {
30            pool: db.pool_arc()?,
31            write_pool: db.write_pool_arc()?,
32        })
33    }
34    pub async fn find_by_id(&self, session_id: &SessionId) -> Result<Option<AnalyticsSession>> {
35        queries::find_by_id(&self.write_pool, session_id).await
36    }
37    pub async fn find_active_by_id(
38        &self,
39        session_id: &SessionId,
40    ) -> Result<Option<ActiveSessionLookup>> {
41        queries::find_active_by_id(&self.write_pool, session_id).await
42    }
43    pub async fn revoke_session(&self, session_id: &SessionId) -> Result<()> {
44        mutations::revoke_session(&self.write_pool, session_id).await
45    }
46    pub async fn revoke_all_for_user(&self, user_id: &UserId) -> Result<u64> {
47        mutations::revoke_all_for_user(&self.write_pool, user_id).await
48    }
49    pub async fn find_by_fingerprint(
50        &self,
51        fingerprint_hash: &str,
52        user_id: &UserId,
53    ) -> Result<Option<AnalyticsSession>> {
54        queries::find_by_fingerprint(&self.pool, fingerprint_hash, user_id).await
55    }
56    pub async fn list_active_by_user(&self, user_id: &UserId) -> Result<Vec<AnalyticsSession>> {
57        queries::list_active_by_user(&self.pool, user_id).await
58    }
59    pub async fn update_activity(&self, session_id: &SessionId) -> Result<()> {
60        mutations::update_activity(&self.write_pool, session_id).await
61    }
62    pub async fn increment_request_count(&self, session_id: &SessionId) -> Result<()> {
63        mutations::increment_request_count(&self.write_pool, session_id).await
64    }
65    pub async fn increment_task_count(&self, session_id: &SessionId) -> Result<()> {
66        mutations::increment_task_count(&self.write_pool, session_id).await
67    }
68    pub async fn increment_message_count(&self, session_id: &SessionId) -> Result<()> {
69        mutations::increment_message_count(&self.write_pool, session_id).await
70    }
71    pub async fn end_session(&self, session_id: &SessionId) -> Result<()> {
72        mutations::end_session(&self.write_pool, session_id).await
73    }
74    pub async fn mark_as_scanner(&self, session_id: &SessionId) -> Result<()> {
75        mutations::mark_as_scanner(&self.write_pool, session_id).await
76    }
77    pub async fn mark_converted(&self, session_id: &SessionId) -> Result<()> {
78        mutations::mark_converted(&self.write_pool, session_id).await
79    }
80    pub async fn mark_as_behavioral_bot(&self, session_id: &SessionId, reason: &str) -> Result<()> {
81        behavioral::mark_as_behavioral_bot(&self.write_pool, session_id, reason).await
82    }
83    pub async fn check_and_mark_behavioral_bot(
84        &self,
85        session_id: &SessionId,
86        request_count_threshold: i32,
87    ) -> Result<bool> {
88        behavioral::check_and_mark_behavioral_bot(
89            &self.write_pool,
90            session_id,
91            request_count_threshold,
92        )
93        .await
94    }
95    pub async fn cleanup_inactive(&self, inactive_hours: i32) -> Result<u64> {
96        mutations::cleanup_inactive(&self.write_pool, inactive_hours).await
97    }
98    pub async fn count_inactive(&self, inactive_hours: i32) -> Result<i64> {
99        queries::count_inactive(&self.pool, inactive_hours).await
100    }
101    pub async fn count_sessions_missing_geo(&self) -> Result<i64> {
102        mutations::count_sessions_missing_geo(&self.pool).await
103    }
104    pub async fn migrate_user_sessions(
105        &self,
106        old_user_id: &UserId,
107        new_user_id: &UserId,
108    ) -> Result<u64> {
109        mutations::migrate_user_sessions(&self.write_pool, old_user_id, new_user_id).await
110    }
111    pub async fn create_session(&self, params: &CreateSessionParams<'_>) -> Result<()> {
112        mutations::create_session(&self.write_pool, params).await
113    }
114    pub async fn find_recent_by_fingerprint(
115        &self,
116        fingerprint_hash: &str,
117        max_age_seconds: i64,
118    ) -> Result<Option<SessionRecord>> {
119        queries::find_recent_by_fingerprint(&self.write_pool, fingerprint_hash, max_age_seconds)
120            .await
121    }
122    pub async fn increment_ai_usage(
123        &self,
124        session_id: &SessionId,
125        tokens: i32,
126        cost_microdollars: i64,
127    ) -> Result<()> {
128        mutations::increment_ai_usage(&self.write_pool, session_id, tokens, cost_microdollars).await
129    }
130    pub async fn update_behavioral_detection(
131        &self,
132        session_id: &SessionId,
133        score: i32,
134        is_behavioral_bot: bool,
135        reason: Option<&str>,
136    ) -> Result<()> {
137        behavioral::update_behavioral_detection(
138            &self.write_pool,
139            session_id,
140            score,
141            is_behavioral_bot,
142            reason,
143        )
144        .await
145    }
146    pub async fn count_sessions_by_fingerprint(
147        &self,
148        fingerprint_hash: &str,
149        window_hours: i64,
150    ) -> Result<i64> {
151        behavioral_queries::count_sessions_by_fingerprint(
152            &self.write_pool,
153            fingerprint_hash,
154            window_hours,
155        )
156        .await
157    }
158    pub async fn get_session_for_behavioral_analysis(
159        &self,
160        session_id: &SessionId,
161    ) -> Result<Option<SessionBehavioralData>> {
162        behavioral_queries::get_session_for_behavioral_analysis(&self.write_pool, session_id).await
163    }
164    pub async fn count_unique_ips_by_fingerprint(
165        &self,
166        fingerprint_hash: &str,
167        window_days: i64,
168    ) -> Result<i64> {
169        behavioral_queries::count_unique_ips_by_fingerprint(
170            &self.write_pool,
171            fingerprint_hash,
172            window_days,
173        )
174        .await
175    }
176    pub async fn get_session_starts_by_fingerprint(
177        &self,
178        fingerprint_hash: &str,
179        window_days: i64,
180    ) -> Result<Vec<DateTime<Utc>>> {
181        behavioral_queries::get_session_starts_by_fingerprint(
182            &self.write_pool,
183            fingerprint_hash,
184            window_days,
185        )
186        .await
187    }
188    pub async fn get_session_velocity(
189        &self,
190        session_id: &SessionId,
191    ) -> Result<(Option<i64>, Option<i64>)> {
192        behavioral_queries::get_session_velocity(&self.write_pool, session_id).await
193    }
194}
195
196mod ai_provider;
197pub use ai_provider::UsersAiSessionProvider;
198
199mod fingerprint;
200
201mod store;