Skip to main content

agentsight_capture/analyzers/
http_decompressor.rs

1// SPDX-License-Identifier: MIT
2// Copyright (c) 2026 eunomia-bpf org.
3
4use 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}