systemprompt_api/services/middleware/analytics/
detection.rs1use chrono::{DateTime, Utc};
7use std::sync::Arc;
8
9use systemprompt_analytics::{
10 BehavioralAnalysisInput, BehavioralBotDetector, SessionSignalsRepository,
11};
12use systemprompt_identifiers::SessionId;
13use systemprompt_traits::{BackgroundTasks, DynSessionStore, SessionStore};
14
15const BEHAVIORAL_FINGERPRINT_WINDOW_DAYS: i64 = 45;
16
17#[derive(Debug, Clone)]
18pub struct DetectionSubject {
19 pub session_id: SessionId,
20 pub fingerprint_hash: Option<String>,
21 pub user_agent: Option<String>,
22 pub request_count: i64,
23}
24
25pub(super) fn spawn_behavioral_detection_task(
26 background: &BackgroundTasks,
27 sessions: DynSessionStore,
28 signals: Arc<SessionSignalsRepository>,
29 subject: DetectionSubject,
30) {
31 background.spawn("analytics_behavioral_detection", async move {
32 let session_id = subject.session_id.clone();
33 let input = collect_analysis_input(&*sessions, &signals, subject).await;
34
35 let result = BehavioralBotDetector::analyze(&input);
36
37 if result.score > 0
38 && let Err(e) = sessions
39 .update_behavioral_detection(
40 &session_id,
41 result.score,
42 result.is_suspicious,
43 result.reason.as_deref(),
44 )
45 .await
46 {
47 tracing::error!(error = %e, "Failed to update behavioral detection");
48 }
49 });
50}
51
52pub async fn collect_analysis_input(
53 sessions: &dyn SessionStore,
54 signals: &SessionSignalsRepository,
55 subject: DetectionSubject,
56) -> BehavioralAnalysisInput {
57 let DetectionSubject {
58 session_id,
59 fingerprint_hash,
60 user_agent,
61 request_count,
62 } = subject;
63 let fingerprint = fingerprint_stats(sessions, signals, fingerprint_hash.as_deref()).await;
64
65 let endpoints_accessed = signals
66 .get_endpoint_sequence(&session_id)
67 .await
68 .unwrap_or_else(|e| {
69 tracing::debug!(error = %e, "Failed to get endpoint sequence");
70 Vec::new()
71 });
72
73 let request_timestamps = signals
74 .get_request_timestamps(&session_id)
75 .await
76 .unwrap_or_else(|e| {
77 tracing::debug!(error = %e, "Failed to get request timestamps");
78 Vec::new()
79 });
80
81 let total_site_pages = signals.get_total_content_pages().await.unwrap_or_else(|e| {
82 tracing::debug!(error = %e, "Failed to get total content pages");
83 100
84 });
85
86 let has_javascript_events = signals
87 .has_analytics_events(&session_id)
88 .await
89 .unwrap_or(false);
90
91 let timeline = session_timeline(sessions, &session_id, request_count).await;
92
93 BehavioralAnalysisInput {
94 session_id,
95 fingerprint_hash,
96 user_agent,
97 request_count: timeline.request_count,
98 started_at: timeline.started_at,
99 last_activity_at: timeline.last_activity_at,
100 endpoints_accessed,
101 total_site_pages,
102 fingerprint_session_count: fingerprint.session_count,
103 fingerprint_unique_ip_count: fingerprint.unique_ip_count,
104 fingerprint_engagement_event_count: fingerprint.engagement_event_count,
105 fingerprint_session_starts: fingerprint.session_starts,
106 request_timestamps,
107 has_javascript_events,
108 landing_page: timeline.landing_page,
109 entry_url: timeline.entry_url,
110 }
111}
112
113struct FingerprintStats {
114 session_count: i64,
115 unique_ip_count: i64,
116 engagement_event_count: i64,
117 session_starts: Vec<DateTime<Utc>>,
118}
119
120async fn fingerprint_stats(
121 sessions: &dyn SessionStore,
122 signals: &SessionSignalsRepository,
123 fingerprint: Option<&str>,
124) -> FingerprintStats {
125 let Some(fp) = fingerprint else {
126 return FingerprintStats {
127 session_count: 1,
128 unique_ip_count: 0,
129 engagement_event_count: 0,
130 session_starts: Vec::new(),
131 };
132 };
133
134 let session_count = sessions
135 .count_sessions_by_fingerprint(fp, 24)
136 .await
137 .unwrap_or_else(|e| {
138 tracing::debug!(error = %e, "Failed to count fingerprint sessions");
139 1
140 });
141
142 let unique_ip_count = sessions
143 .count_unique_ips_by_fingerprint(fp, BEHAVIORAL_FINGERPRINT_WINDOW_DAYS)
144 .await
145 .unwrap_or_else(|e| {
146 tracing::debug!(error = %e, "Failed to count fingerprint unique IPs");
147 0
148 });
149
150 let engagement_event_count = signals
151 .count_engagement_events_by_fingerprint(fp, BEHAVIORAL_FINGERPRINT_WINDOW_DAYS)
152 .await
153 .unwrap_or_else(|e| {
154 tracing::debug!(error = %e, "Failed to count fingerprint engagement events");
155 0
156 });
157
158 let session_starts = sessions
159 .get_session_starts_by_fingerprint(fp, BEHAVIORAL_FINGERPRINT_WINDOW_DAYS)
160 .await
161 .unwrap_or_else(|e| {
162 tracing::debug!(error = %e, "Failed to load fingerprint session starts");
163 Vec::new()
164 });
165
166 FingerprintStats {
167 session_count,
168 unique_ip_count,
169 engagement_event_count,
170 session_starts,
171 }
172}
173
174struct SessionTimeline {
175 started_at: DateTime<Utc>,
176 last_activity_at: DateTime<Utc>,
177 request_count: i64,
178 landing_page: Option<String>,
179 entry_url: Option<String>,
180}
181
182async fn session_timeline(
183 sessions: &dyn SessionStore,
184 session_id: &SessionId,
185 request_count: i64,
186) -> SessionTimeline {
187 let session_data = sessions
188 .get_session_for_behavioral_analysis(session_id)
189 .await
190 .map_err(|e| {
191 tracing::debug!(error = %e, "Failed to get session for behavioral analysis");
192 e
193 })
194 .ok()
195 .flatten();
196
197 session_data.map_or_else(
198 || {
199 let now = Utc::now();
200 SessionTimeline {
201 started_at: now,
202 last_activity_at: now,
203 request_count,
204 landing_page: None,
205 entry_url: None,
206 }
207 },
208 |s| SessionTimeline {
209 started_at: s.started_at,
210 last_activity_at: s.last_activity_at,
211 request_count: s.request_count.map_or(request_count, i64::from),
212 landing_page: s.landing_page,
213 entry_url: s.entry_url,
214 },
215 )
216}