agentsight_capture/analyzers/
auth_header_remover.rs1use super::{Analyzer, AnalyzerError};
5use crate::runners::EventStream;
6use async_trait::async_trait;
7use futures::stream::StreamExt;
8use serde_json::Value;
9
10#[derive(Debug)]
13pub struct AuthHeaderRemover {
14 auth_headers: Vec<String>,
16 debug: bool,
18}
19
20impl AuthHeaderRemover {
21 pub fn new() -> Self {
23 Self {
24 auth_headers: vec![
25 "authorization".to_string(),
26 "x-api-key".to_string(),
27 "x-auth-token".to_string(),
28 "bearer".to_string(),
29 "token".to_string(),
30 "x-access-token".to_string(),
31 "x-session-token".to_string(),
32 "cookie".to_string(),
33 "set-cookie".to_string(),
34 ],
35 debug: false,
36 }
37 }
38
39 fn remove_auth_headers(&self, mut event_data: Value) -> Value {
41 if event_data
43 .get("message_type")
44 .and_then(|v| v.as_str())
45 .is_none()
46 {
47 return event_data;
48 }
49
50 let mut headers_removed = Vec::new();
51
52 if let Some(headers_obj) = event_data
54 .get_mut("headers")
55 .and_then(|v| v.as_object_mut())
56 {
57 let keys_to_remove: Vec<String> = headers_obj
59 .keys()
60 .filter(|key| {
61 self.auth_headers
62 .iter()
63 .any(|auth_header| key.to_lowercase() == auth_header.to_lowercase())
64 })
65 .cloned()
66 .collect();
67
68 for key in keys_to_remove {
70 if headers_obj.remove(&key).is_some() {
71 headers_removed.push(key);
72 }
73 }
74 }
75
76 if self.debug && !headers_removed.is_empty() {
78 eprintln!(
79 "[AuthHeaderRemover DEBUG] Removed headers: {:?}",
80 headers_removed
81 );
82 }
83
84 event_data
85 }
86}
87
88impl Default for AuthHeaderRemover {
89 fn default() -> Self {
90 Self::new()
91 }
92}
93
94#[async_trait]
95impl Analyzer for AuthHeaderRemover {
96 async fn process(&mut self, stream: EventStream) -> Result<EventStream, AnalyzerError> {
97 let auth_headers = self.auth_headers.clone();
98 let debug = self.debug;
99
100 let processed_stream = stream.map(move |mut event| {
101 if event.source == "http_parser" {
103 event.data = AuthHeaderRemover {
104 auth_headers: auth_headers.clone(),
105 debug,
106 }
107 .remove_auth_headers(event.data);
108 }
109 event
110 });
111
112 Ok(Box::pin(processed_stream))
113 }
114}
115
116#[cfg(test)]
117mod tests {
118 use super::*;
119 use crate::event::Event;
120 use futures::stream;
121 use serde_json::json;
122
123 #[tokio::test]
124 async fn test_auth_header_removal() {
125 let mut analyzer = AuthHeaderRemover::new();
126
127 let event_data = json!({
128 "message_type": "request",
129 "method": "GET",
130 "path": "/api/test",
131 "headers": {
132 "authorization": "Bearer token123",
133 "content-type": "application/json",
134 "x-api-key": "secret-key",
135 "user-agent": "test-client"
136 }
137 });
138
139 let test_event = Event::new(
140 "http_parser".to_string(),
141 1234,
142 "http_parser".to_string(),
143 event_data,
144 );
145 let events = vec![test_event];
146
147 let input_stream: EventStream = Box::pin(stream::iter(events));
148 let output_stream = analyzer.process(input_stream).await.unwrap();
149
150 let collected: Vec<_> = output_stream.collect().await;
151
152 assert_eq!(collected.len(), 1);
153 let headers = collected[0]
154 .data
155 .get("headers")
156 .and_then(|v| v.as_object())
157 .unwrap();
158
159 assert!(!headers.contains_key("authorization"));
161 assert!(!headers.contains_key("x-api-key"));
162
163 assert!(headers.contains_key("content-type"));
165 assert!(headers.contains_key("user-agent"));
166 }
167
168 #[tokio::test]
169 async fn test_non_http_events_passthrough() {
170 let mut analyzer = AuthHeaderRemover::new();
171
172 let event_data = json!({
173 "type": "process",
174 "pid": 1234,
175 "command": "test"
176 });
177
178 let test_event = Event::new(
179 "process".to_string(),
180 1234,
181 "process".to_string(),
182 event_data.clone(),
183 );
184 let events = vec![test_event];
185
186 let input_stream: EventStream = Box::pin(stream::iter(events));
187 let output_stream = analyzer.process(input_stream).await.unwrap();
188
189 let collected: Vec<_> = output_stream.collect().await;
190
191 assert_eq!(collected.len(), 1);
192 assert_eq!(collected[0].data, event_data);
193 }
194
195 #[tokio::test]
196 async fn test_case_insensitive_matching() {
197 let mut analyzer = AuthHeaderRemover::new();
198
199 let event_data = json!({
200 "message_type": "request",
201 "headers": {
202 "Authorization": "Bearer token123",
203 "X-API-KEY": "secret-key",
204 "Content-Type": "application/json"
205 }
206 });
207
208 let test_event = Event::new(
209 "http_parser".to_string(),
210 1234,
211 "http_parser".to_string(),
212 event_data,
213 );
214 let events = vec![test_event];
215
216 let input_stream: EventStream = Box::pin(stream::iter(events));
217 let output_stream = analyzer.process(input_stream).await.unwrap();
218
219 let collected: Vec<_> = output_stream.collect().await;
220
221 let headers = collected[0]
222 .data
223 .get("headers")
224 .and_then(|v| v.as_object())
225 .unwrap();
226
227 assert!(!headers.contains_key("Authorization"));
229 assert!(!headers.contains_key("X-API-KEY"));
230
231 assert!(headers.contains_key("Content-Type"));
233 }
234
235 #[tokio::test]
236 async fn test_no_headers_field() {
237 let mut analyzer = AuthHeaderRemover::new();
238
239 let event_data = json!({
240 "message_type": "request",
241 "method": "GET",
242 "path": "/api/test"
243 });
244
245 let test_event = Event::new(
246 "http_parser".to_string(),
247 1234,
248 "http_parser".to_string(),
249 event_data.clone(),
250 );
251 let events = vec![test_event];
252
253 let input_stream: EventStream = Box::pin(stream::iter(events));
254 let output_stream = analyzer.process(input_stream).await.unwrap();
255
256 let collected: Vec<_> = output_stream.collect().await;
257
258 assert_eq!(collected.len(), 1);
260 assert_eq!(collected[0].data, event_data);
261 }
262}