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