agentsight_capture/analyzers/
http_decompressor.rs1use super::{Analyzer, AnalyzerError};
5use crate::event::Event;
6use crate::runners::EventStream;
7use async_trait::async_trait;
8use flate2::read::{DeflateDecoder, GzDecoder, ZlibDecoder};
9use futures::stream::StreamExt;
10use serde_json::{Value, json};
11use std::io::{Cursor, Read};
12
13#[derive(Default)]
14pub struct HTTPDecompressor;
15
16impl HTTPDecompressor {
17 pub fn new() -> Self {
18 Self
19 }
20
21 fn process_event(mut event: Event) -> Event {
22 if event.source != "http_parser" {
23 return event;
24 }
25 if event.data.get("message_type").and_then(Value::as_str) != Some("response") {
26 return event;
27 }
28
29 let Some(encoding) = content_encoding(&event.data) else {
30 return event;
31 };
32 let Some(compressed) = http_body_bytes(&event.data) else {
33 return event;
34 };
35 let payload = if is_chunked(&event.data) {
36 decode_chunked(&compressed).unwrap_or(compressed)
37 } else {
38 compressed
39 };
40
41 let Ok(decompressed) = decompress_body(&encoding, &payload) else {
42 return event;
43 };
44 let decompressed_len = decompressed.len();
45 let decompressed_body = String::from_utf8_lossy(&decompressed).to_string();
46
47 event.data["body"] = Value::String(decompressed_body);
48 event.data["body_hex"] = Value::String(hex::encode(&decompressed));
49 event.data["has_body"] = Value::Bool(decompressed_len > 0);
50 event.data["content_length"] = json!(decompressed_len);
51 event.data["decompressed"] = Value::Bool(true);
52 event.data["original_content_encoding"] = Value::String(encoding);
53 event.data["decompressed_body_size"] = json!(decompressed_len);
54
55 if let Some(headers) = event.data.get_mut("headers").and_then(Value::as_object_mut) {
56 headers.remove("content-encoding");
57 headers.remove("Content-Encoding");
58 headers.remove("content-length");
59 headers.remove("Content-Length");
60 }
61
62 event
63 }
64}
65
66#[async_trait]
67impl Analyzer for HTTPDecompressor {
68 async fn process(&mut self, stream: EventStream) -> Result<EventStream, AnalyzerError> {
69 let processed = stream.map(Self::process_event);
70 Ok(Box::pin(processed))
71 }
72}
73
74fn content_encoding(data: &Value) -> Option<String> {
75 let headers = data.get("headers")?.as_object()?;
76 headers
77 .iter()
78 .find(|(key, _)| key.eq_ignore_ascii_case("content-encoding"))
79 .and_then(|(_, value)| value.as_str())
80 .map(|value| {
81 value
82 .split(',')
83 .next()
84 .unwrap_or(value)
85 .trim()
86 .to_ascii_lowercase()
87 })
88 .filter(|encoding| {
89 matches!(
90 encoding.as_str(),
91 "gzip" | "x-gzip" | "deflate" | "br" | "brotli" | "zstd" | "zstandard"
92 )
93 })
94}
95
96fn is_chunked(data: &Value) -> bool {
97 if data
98 .get("is_chunked")
99 .and_then(Value::as_bool)
100 .unwrap_or(false)
101 {
102 return true;
103 }
104 let Some(headers) = data.get("headers").and_then(Value::as_object) else {
105 return false;
106 };
107 headers.iter().any(|(key, value)| {
108 key.eq_ignore_ascii_case("transfer-encoding")
109 && value
110 .as_str()
111 .is_some_and(|v| v.to_ascii_lowercase().contains("chunked"))
112 })
113}
114
115fn decode_chunked(bytes: &[u8]) -> Option<Vec<u8>> {
116 let cursor = Cursor::new(bytes);
117 let mut decoder = chunked_transfer::Decoder::new(cursor);
118 let mut out = Vec::new();
119 decoder.read_to_end(&mut out).ok()?;
120 Some(out)
121}
122
123fn http_body_bytes(data: &Value) -> Option<Vec<u8>> {
124 data.get("body_hex")
125 .and_then(Value::as_str)
126 .and_then(|value| hex::decode(value).ok())
127 .or_else(|| {
128 data.get("body")
129 .and_then(Value::as_str)
130 .map(http_body_string_to_bytes)
131 })
132}
133
134fn decompress_body(encoding: &str, body: &[u8]) -> Result<Vec<u8>, std::io::Error> {
135 match encoding {
136 "gzip" | "x-gzip" => read_all(GzDecoder::new(Cursor::new(body))),
137 "deflate" => read_all(ZlibDecoder::new(Cursor::new(body)))
138 .or_else(|_| read_all(DeflateDecoder::new(Cursor::new(body)))),
139 "br" | "brotli" => read_all(brotli::Decompressor::new(Cursor::new(body), 4096)),
140 "zstd" | "zstandard" => {
141 zstd::stream::decode_all(Cursor::new(body)).map_err(std::io::Error::other)
142 }
143 _ => Ok(body.to_vec()),
144 }
145}
146
147fn read_all<R: Read>(mut reader: R) -> Result<Vec<u8>, std::io::Error> {
148 let mut out = Vec::new();
149 reader.read_to_end(&mut out)?;
150 Ok(out)
151}
152
153fn http_body_string_to_bytes(data: &str) -> Vec<u8> {
154 let mut bytes = Vec::with_capacity(data.len());
155 for ch in data.chars() {
156 let code = ch as u32;
157 if code <= 0xff {
158 bytes.push(code as u8);
159 } else {
160 let mut buf = [0u8; 4];
161 bytes.extend_from_slice(ch.encode_utf8(&mut buf).as_bytes());
162 }
163 }
164 bytes
165}
166
167#[cfg(test)]
168mod tests {
169 use super::*;
170 use flate2::Compression;
171 use flate2::write::{DeflateEncoder, GzEncoder, ZlibEncoder};
172 use std::io::Write;
173
174 fn http_response(body: Vec<u8>, encoding: &str) -> Event {
175 Event::new(
176 "http_parser".to_string(),
177 123,
178 "agent".to_string(),
179 json!({
180 "tid": 7,
181 "message_type": "response",
182 "status_code": 200,
183 "headers": {
184 "content-encoding": encoding,
185 "content-type": "text/event-stream"
186 },
187 "body": bytes_to_http_body_string(&body),
188 "has_body": true,
189 "is_chunked": false
190 }),
191 )
192 }
193
194 fn bytes_to_http_body_string(bytes: &[u8]) -> String {
195 bytes.iter().map(|b| char::from(*b)).collect()
196 }
197
198 #[test]
199 fn decompresses_gzip_response_body() {
200 let mut encoder = GzEncoder::new(Vec::new(), Compression::default());
201 encoder
202 .write_all(b"data: {\"usage\":{\"input_tokens\":1}}\n\n")
203 .unwrap();
204 let event =
205 HTTPDecompressor::process_event(http_response(encoder.finish().unwrap(), "gzip"));
206
207 assert_eq!(event.data["decompressed"], true);
208 assert_eq!(
209 event.data["body"].as_str().unwrap(),
210 "data: {\"usage\":{\"input_tokens\":1}}\n\n"
211 );
212 assert!(event.data["headers"].get("content-encoding").is_none());
213 }
214
215 #[test]
216 fn decompresses_zlib_and_raw_deflate() {
217 let mut zlib = ZlibEncoder::new(Vec::new(), Compression::default());
218 zlib.write_all(b"zlib-body").unwrap();
219 let event =
220 HTTPDecompressor::process_event(http_response(zlib.finish().unwrap(), "deflate"));
221 assert_eq!(event.data["body"].as_str().unwrap(), "zlib-body");
222
223 let mut raw = DeflateEncoder::new(Vec::new(), Compression::default());
224 raw.write_all(b"raw-body").unwrap();
225 let event =
226 HTTPDecompressor::process_event(http_response(raw.finish().unwrap(), "deflate"));
227 assert_eq!(event.data["body"].as_str().unwrap(), "raw-body");
228 }
229
230 #[test]
231 fn decompresses_brotli_and_zstd_response_bodies() {
232 let mut br = Vec::new();
233 {
234 let mut encoder = brotli::CompressorWriter::new(&mut br, 4096, 5, 22);
235 encoder.write_all(b"brotli-body").unwrap();
236 }
237 let event = HTTPDecompressor::process_event(http_response(br, "br"));
238 assert_eq!(event.data["body"].as_str().unwrap(), "brotli-body");
239
240 let zstd = zstd::stream::encode_all(Cursor::new(b"zstd-body"), 1).unwrap();
241 let event = HTTPDecompressor::process_event(http_response(zstd, "zstd"));
242 assert_eq!(event.data["body"].as_str().unwrap(), "zstd-body");
243 }
244
245 #[test]
246 fn leaves_requests_and_unknown_encodings_unchanged() {
247 let mut request = http_response(b"abc".to_vec(), "gzip");
248 request.data["message_type"] = Value::String("request".to_string());
249 assert_eq!(
250 HTTPDecompressor::process_event(request).data["body"].as_str(),
251 Some("abc")
252 );
253
254 let event = http_response(b"abc".to_vec(), "identity");
255 assert!(
256 HTTPDecompressor::process_event(event)
257 .data
258 .get("decompressed")
259 .is_none()
260 );
261 }
262
263 #[test]
264 fn prefers_body_hex_for_non_utf8_compressed_payloads() {
265 let mut encoder = GzEncoder::new(Vec::new(), Compression::default());
266 encoder
267 .write_all(b"data: {\"usage\":{\"total_tokens\":7}}\n\n")
268 .unwrap();
269 let compressed = encoder.finish().unwrap();
270 let mut event = http_response(Vec::new(), "gzip");
271 event.data["body"] = Value::String("\u{fffd}\u{fffd}".to_string());
272 event.data["body_hex"] = Value::String(hex::encode(compressed));
273
274 let event = HTTPDecompressor::process_event(event);
275
276 assert_eq!(
277 event.data["body"].as_str().unwrap(),
278 "data: {\"usage\":{\"total_tokens\":7}}\n\n"
279 );
280 assert_eq!(
281 hex::decode(event.data["body_hex"].as_str().unwrap()).unwrap(),
282 b"data: {\"usage\":{\"total_tokens\":7}}\n\n"
283 );
284 }
285}