systemprompt_api/services/middleware/analytics/
events.rs1use serde_json::json;
14use std::sync::Arc;
15
16use systemprompt_logging::{AnalyticsEvent, AnalyticsRepository};
17use systemprompt_models::routing::EventMetadata;
18use systemprompt_models::{RequestContext, RouteClassifier};
19use systemprompt_traits::BackgroundTasks;
20
21#[derive(Debug)]
22pub struct AnalyticsEventParams {
23 pub req_ctx: RequestContext,
24 pub endpoint: String,
25 pub path: String,
26 pub method: String,
27 pub uri: http::Uri,
28 pub status_code: u16,
29 pub response_time_ms: u64,
30 pub user_agent: Option<String>,
31 pub referer: Option<String>,
32 pub html_response: bool,
33}
34
35#[must_use]
36pub fn event_metadata_for(classified: EventMetadata, html_response: bool) -> EventMetadata {
37 if classified == EventMetadata::HTML_CONTENT && !html_response {
38 EventMetadata::API_REQUEST
39 } else {
40 classified
41 }
42}
43
44pub(super) fn spawn_analytics_event_task(
45 background: &BackgroundTasks,
46 analytics_repo: Arc<AnalyticsRepository>,
47 route_classifier: Arc<RouteClassifier>,
48 params: AnalyticsEventParams,
49) {
50 let sanitized_uri = sanitize_uri(¶ms.uri);
51
52 background.spawn("analytics_event", async move {
53 let message = format!(
54 "HTTP {} - {} {}",
55 params.status_code, params.method, sanitized_uri
56 );
57 let metadata = json!({
58 "status_code": params.status_code,
59 "method": params.method,
60 "uri": sanitized_uri,
61 "endpoint": params.endpoint,
62 "trace_id": params.req_ctx.trace_id(),
63 "user_agent": params.user_agent,
64 "referer": params.referer
65 });
66
67 let event_metadata = event_metadata_for(
68 route_classifier.get_event_metadata(¶ms.path, ¶ms.method),
69 params.html_response,
70 );
71
72 let severity = if params.status_code >= 500 {
73 "error"
74 } else if params.status_code >= 400 {
75 "warning"
76 } else {
77 "info"
78 };
79
80 let event = AnalyticsEvent {
81 user_id: params.req_ctx.auth.actor.user_id.clone(),
82 session_id: params.req_ctx.request.session_id.clone(),
83 context_id: params.req_ctx.execution.context_id.clone(),
84 event_type: event_metadata.event_type.to_owned(),
85 event_category: event_metadata.event_category.to_owned(),
86 severity: severity.to_owned(),
87 endpoint: Some(params.endpoint),
88 error_code: if params.status_code >= 400 {
89 Some(i32::from(params.status_code))
90 } else {
91 None
92 },
93 response_time_ms: Some(params.response_time_ms as i32),
94 agent_id: None,
95 task_id: params.req_ctx.task_id().cloned(),
96 message: Some(message.clone()),
97 metadata: metadata.clone(),
98 };
99
100 if let Err(e) = analytics_repo.log_event(&event).await {
101 tracing::error!(error = %e, "Failed to log analytics event");
102 }
103
104 if params.status_code >= 500 {
105 tracing::error!(module = event_metadata.log_module, message = %message, metadata = ?metadata, "HTTP error");
106 }
107 });
108}
109
110pub fn sanitize_uri(uri: &http::Uri) -> String {
111 let path = uri.path();
112
113 uri.query().map_or_else(
114 || path.to_owned(),
115 |query| {
116 let sanitized_params: Vec<String> = query
117 .split('&')
118 .map(|param| {
119 param.split_once('=').map_or_else(
120 || param.to_owned(),
121 |(key, value)| {
122 let key_lower = key.to_lowercase();
123 if is_sensitive_key(&key_lower) {
124 format!("{key}=[REDACTED]")
125 } else {
126 format!("{key}={value}")
127 }
128 },
129 )
130 })
131 .collect();
132
133 format!("{path}?{}", sanitized_params.join("&"))
134 },
135 )
136}
137
138pub fn is_sensitive_key(key: &str) -> bool {
139 matches!(
140 key,
141 "token" | "password" | "api_key" | "apikey" | "secret" | "authorization" | "auth"
142 )
143}