systemprompt_users/sessions/
mod.rs1mod 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;