Skip to main content

systemprompt_api/services/middleware/analytics/
detection.rs

1//! Behavioral bot detection input collection over fingerprint history.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use 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}