Skip to main content

agentsight_capture/analyzers/
http_parser.rs

1// SPDX-License-Identifier: MIT
2// Copyright (c) 2026 eunomia-bpf org.
3
4use super::protocol_events::HTTPEvent;
5use super::{Analyzer, AnalyzerError};
6use crate::event::Event;
7use crate::runners::EventStream;
8use async_trait::async_trait;
9use flate2::{Decompress, FlushDecompress};
10use futures::{stream, stream::StreamExt};
11use hpack::Decoder as HpackDecoder;
12use std::collections::HashMap;
13
14const MAX_HTTP2_STREAMS: usize = 1024;
15const MAX_HTTP2_PENDING_HEADERS: usize = 1024;
16const MAX_HTTP_BODY_BYTES: usize = 1024 * 1024;
17const MAX_HTTP2_HEADER_BLOCK_BYTES: usize = 64 * 1024;
18
19/// HTTP Parser Analyzer that parses SSL traffic into HTTP requests/responses
20pub struct HTTPParser {
21    /// Flag to include raw data in parsed events (default: true)
22    include_raw_data: bool,
23    http2: HTTP2State,
24    websocket: WebSocketState,
25}
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28enum HTTP2Direction {
29    Request,
30    Response,
31}
32
33#[derive(Default)]
34struct HTTP2StreamState {
35    request_headers: HashMap<String, String>,
36    response_headers: HashMap<String, String>,
37    request_body: Vec<u8>,
38    response_body: Vec<u8>,
39    request_emitted: bool,
40    response_emitted: bool,
41}
42
43struct PendingHTTP2Headers {
44    direction: HTTP2Direction,
45    block: Vec<u8>,
46}
47
48struct HTTP2Frame<'a> {
49    frame_type: u8,
50    flags: u8,
51    stream_id: u32,
52    payload: &'a [u8],
53}
54
55struct HTTP2State {
56    request_decoder: HpackDecoder<'static>,
57    response_decoder: HpackDecoder<'static>,
58    streams: HashMap<(u64, u32), HTTP2StreamState>,
59    pending_headers: HashMap<(u64, u32), PendingHTTP2Headers>,
60}
61
62#[derive(Default)]
63struct WebSocketState {
64    connections: HashMap<u32, WebSocketConnection>,
65}
66
67struct WebSocketConnection {
68    path: String,
69    headers: HashMap<String, String>,
70    inflater: Decompress,
71}
72
73impl Default for HTTP2State {
74    fn default() -> Self {
75        Self {
76            request_decoder: HpackDecoder::new(),
77            response_decoder: HpackDecoder::new(),
78            streams: HashMap::new(),
79            pending_headers: HashMap::new(),
80        }
81    }
82}
83
84#[derive(Clone, PartialEq, Debug)]
85pub enum HTTPMessageType {
86    Request,
87    Response,
88}
89
90/// Parsed HTTP message
91#[derive(Clone, Debug)]
92pub struct HTTPMessage {
93    pub message_type: HTTPMessageType,
94    pub first_line: String,
95    pub headers: HashMap<String, String>,
96    pub body: Option<String>,
97    pub raw_data: String,
98    // Request-specific fields
99    pub method: Option<String>,
100    pub path: Option<String>,
101    pub protocol: Option<String>,
102    // Response-specific fields
103    pub status_code: Option<u16>,
104    pub status_text: Option<String>,
105}
106
107impl Default for HTTPParser {
108    fn default() -> Self {
109        Self::new()
110    }
111}
112
113impl HTTPParser {
114    /// Create a new HTTPParser with default settings (raw data included)
115    pub fn new() -> Self {
116        HTTPParser {
117            include_raw_data: true,
118            http2: HTTP2State::default(),
119            websocket: WebSocketState::default(),
120        }
121    }
122
123    /// Disable raw data inclusion
124    pub fn disable_raw_data(mut self) -> Self {
125        self.include_raw_data = false;
126        self
127    }
128
129    /// Check if SSL data contains HTTP protocol data
130    pub fn is_http_data(data: &str) -> bool {
131        // Look for HTTP patterns
132        let has_http_request = data.contains("HTTP/1.")
133            && (data.contains("GET ")
134                || data.contains("POST ")
135                || data.contains("PUT ")
136                || data.contains("DELETE ")
137                || data.contains("HEAD ")
138                || data.contains("OPTIONS ")
139                || data.contains("PATCH "));
140
141        let has_http_response = data.starts_with("HTTP/1.") || data.contains("\r\nHTTP/1.");
142
143        // Look for common HTTP headers
144        let has_http_headers = data.contains("Content-Type:")
145            || data.contains("content-type:")
146            || data.contains("Host:")
147            || data.contains("host:")
148            || data.contains("User-Agent:")
149            || data.contains("user-agent:");
150
151        has_http_request || has_http_response || has_http_headers
152    }
153
154    /// Parse HTTP message from accumulated data
155    pub fn parse_http_message(data: &str) -> Option<HTTPMessage> {
156        let lines: Vec<&str> = data.split("\r\n").collect();
157
158        if lines.is_empty() {
159            return None;
160        }
161
162        let first_line = lines[0];
163        let mut headers = HashMap::new();
164        let mut body_start = None;
165        let mut message_type = HTTPMessageType::Request;
166        let mut method = None;
167        let mut path = None;
168        let mut protocol = None;
169        let mut status_code = None;
170        let mut status_text = None;
171
172        // Parse first line to determine message type
173        if first_line.starts_with("HTTP/") {
174            // Response
175            message_type = HTTPMessageType::Response;
176            let parts: Vec<&str> = first_line.splitn(3, ' ').collect();
177            if parts.len() >= 2 {
178                if let Ok(code) = parts[1].parse::<u16>() {
179                    status_code = Some(code);
180                }
181                if parts.len() >= 3 {
182                    status_text = Some(parts[2].to_string());
183                }
184                protocol = Some(parts[0].to_string());
185            }
186        } else {
187            // Request
188            let parts: Vec<&str> = first_line.splitn(3, ' ').collect();
189            if parts.len() < 3
190                || !matches!(
191                    parts[0],
192                    "GET" | "POST" | "PUT" | "DELETE" | "HEAD" | "OPTIONS" | "PATCH"
193                )
194                || !parts[2].starts_with("HTTP/")
195            {
196                return None;
197            }
198            method = Some(parts[0].to_string());
199            path = Some(parts[1].to_string());
200            protocol = Some(parts[2].to_string());
201        }
202
203        // Parse headers
204        for (i, line) in lines.iter().enumerate().skip(1) {
205            if line.is_empty() {
206                body_start = Some(i + 1);
207                break;
208            }
209            if let Some(colon_pos) = line.find(':') {
210                let key = line[..colon_pos].trim().to_lowercase();
211                let value = line[colon_pos + 1..].trim().to_string();
212                headers.insert(key, value);
213            }
214        }
215
216        // Extract body if present
217        let body = if let Some(start) = body_start {
218            if start < lines.len() {
219                let body_lines: Vec<&str> = lines[start..].to_vec();
220                let body_content = body_lines.join("\r\n");
221                if !body_content.trim().is_empty() {
222                    Some(body_content)
223                } else {
224                    None
225                }
226            } else {
227                None
228            }
229        } else {
230            None
231        };
232
233        Some(HTTPMessage {
234            message_type,
235            first_line: first_line.to_string(),
236            headers,
237            body,
238            raw_data: data.to_string(),
239            method,
240            path,
241            protocol,
242            status_code,
243            status_text,
244        })
245    }
246
247    /// Create HTTP event from parsed message
248    fn create_http_event(
249        tid: u64,
250        parsed_message: HTTPMessage,
251        original_event: &Event,
252        include_raw_data: bool,
253    ) -> Event {
254        let message_type_str = match parsed_message.message_type {
255            HTTPMessageType::Request => "request",
256            HTTPMessageType::Response => "response",
257        };
258
259        // Determine content properties
260        let content_length = parsed_message
261            .headers
262            .get("content-length")
263            .and_then(|v| v.parse::<usize>().ok());
264        let is_chunked = parsed_message
265            .headers
266            .get("transfer-encoding")
267            .map(|v| v.to_lowercase().contains("chunked"))
268            .unwrap_or(false);
269        let has_body = parsed_message.body.is_some();
270        let body_hex = parsed_message
271            .body
272            .as_deref()
273            .map(ssl_json_string_to_bytes)
274            .map(hex::encode);
275
276        // Calculate total size from parsed components
277        let total_size = parsed_message.first_line.len() +
278            parsed_message.headers.iter().map(|(k, v)| k.len() + v.len() + 4).sum::<usize>() + // +4 for ": \r\n"
279            parsed_message.body.as_ref().map(|b| b.len()).unwrap_or(0) +
280            4; // +4 for \r\n\r\n separator
281
282        HTTPEvent {
283            tid,
284            message_type: message_type_str.to_string(),
285            first_line: parsed_message.first_line,
286            method: parsed_message.method,
287            path: parsed_message.path,
288            protocol: parsed_message.protocol,
289            status_code: parsed_message.status_code,
290            status_text: parsed_message.status_text,
291            headers: parsed_message.headers,
292            body: parsed_message.body,
293            body_hex,
294            total_size,
295            has_body,
296            is_chunked,
297            content_length,
298            original_source: "ssl".to_string(),
299            raw_data: include_raw_data.then_some(parsed_message.raw_data),
300        }
301        .to_event(original_event)
302    }
303
304    /// Handle SSL events (HTTP request/response data)
305    fn handle_ssl_event(
306        http2: &mut HTTP2State,
307        websocket: &mut WebSocketState,
308        event: Event,
309        include_raw_data: bool,
310    ) -> Vec<Event> {
311        let ssl_data = &event.data;
312
313        let data_str = match ssl_data.get("data").and_then(|v| v.as_str()) {
314            Some(s) => s,
315            None => return vec![event],
316        };
317
318        // Only process if it's HTTP data AND can be parsed as a complete HTTP message
319        if Self::is_http_data(data_str)
320            && let Some(parsed_message) = Self::parse_http_message(data_str)
321        {
322            websocket.observe_handshake(&event, &parsed_message);
323            let tid = ssl_data.get("tid").and_then(|v| v.as_u64()).unwrap_or(0);
324            return vec![Self::create_http_event(
325                tid,
326                parsed_message,
327                &event,
328                include_raw_data,
329            )];
330        }
331
332        let data_bytes = ssl_data
333            .get("data_hex")
334            .and_then(|v| v.as_str())
335            .and_then(|v| hex::decode(v).ok())
336            .unwrap_or_else(|| ssl_json_string_to_bytes(data_str));
337        if let Some(events) = websocket.handle_event(&event, &data_bytes, include_raw_data) {
338            return events;
339        }
340        if let Some(events) = http2.handle_event(&event, &data_bytes, include_raw_data) {
341            return events;
342        }
343
344        // If not parseable as HTTP, pass through original event
345        vec![event]
346    }
347}
348
349impl WebSocketState {
350    fn observe_handshake(&mut self, event: &Event, message: &HTTPMessage) {
351        if message.message_type != HTTPMessageType::Request
352            || !message
353                .headers
354                .get("upgrade")
355                .is_some_and(|v| v.eq_ignore_ascii_case("websocket"))
356        {
357            return;
358        }
359        let Some(path) = message.path.as_ref() else {
360            return;
361        };
362        if !path.contains("/v1/responses") && !path.contains("/codex/responses") {
363            return;
364        }
365        self.connections.insert(
366            event.pid,
367            WebSocketConnection {
368                path: path.clone(),
369                headers: message.headers.clone(),
370                inflater: Decompress::new(false),
371            },
372        );
373    }
374
375    fn handle_event(
376        &mut self,
377        event: &Event,
378        bytes: &[u8],
379        include_raw_data: bool,
380    ) -> Option<Vec<Event>> {
381        let connection = self.connections.get_mut(&event.pid)?;
382        let (compressed, mut payload) = parse_masked_websocket_frame(bytes)?;
383        if compressed {
384            payload.extend_from_slice(&[0, 0, 0xff, 0xff]);
385            let mut decoded = Vec::with_capacity(MAX_HTTP_BODY_BYTES);
386            let input_before = connection.inflater.total_in();
387            connection
388                .inflater
389                .decompress_vec(&payload, &mut decoded, FlushDecompress::Sync)
390                .ok()?;
391            if connection.inflater.total_in() - input_before != payload.len() as u64 {
392                return None;
393            }
394            payload = decoded;
395        }
396        let body = String::from_utf8(payload).ok()?;
397        let json: serde_json::Value = serde_json::from_str(&body).ok()?;
398        if json.get("type").and_then(|v| v.as_str()) != Some("response.create") {
399            return None;
400        }
401        Some(vec![create_websocket_request_event(
402            event,
403            &connection.path,
404            &connection.headers,
405            body,
406            include_raw_data,
407        )])
408    }
409}
410
411fn parse_masked_websocket_frame(bytes: &[u8]) -> Option<(bool, Vec<u8>)> {
412    if bytes.len() < 2
413        || bytes[0] & 0x80 == 0
414        || bytes[0] & 0x30 != 0
415        || !matches!(bytes[0] & 0x0f, 1 | 2)
416        || bytes[1] & 0x80 == 0
417    {
418        return None;
419    }
420    let mut offset = 2;
421    let mut payload_len = (bytes[1] & 0x7f) as usize;
422    if payload_len == 126 {
423        payload_len = usize::from(u16::from_be_bytes(
424            bytes.get(offset..offset + 2)?.try_into().ok()?,
425        ));
426        offset += 2;
427    } else if payload_len == 127 {
428        payload_len = usize::try_from(u64::from_be_bytes(
429            bytes.get(offset..offset + 8)?.try_into().ok()?,
430        ))
431        .ok()?;
432        offset += 8;
433    }
434    let mask: [u8; 4] = bytes.get(offset..offset + 4)?.try_into().ok()?;
435    offset += 4;
436    let payload = bytes.get(offset..offset.checked_add(payload_len)?)?;
437    if offset + payload_len != bytes.len() {
438        return None;
439    }
440    Some((
441        bytes[0] & 0x40 != 0,
442        payload
443            .iter()
444            .enumerate()
445            .map(|(i, byte)| byte ^ mask[i % 4])
446            .collect(),
447    ))
448}
449
450fn create_websocket_request_event(
451    original_event: &Event,
452    path: &str,
453    headers: &HashMap<String, String>,
454    body: String,
455    include_raw_data: bool,
456) -> Event {
457    let tid = original_event
458        .data
459        .get("tid")
460        .and_then(|v| v.as_u64())
461        .unwrap_or(0);
462    HTTPEvent {
463        tid,
464        message_type: "request".to_string(),
465        first_line: format!("POST {path} WebSocket"),
466        method: Some("POST".to_string()),
467        path: Some(path.to_string()),
468        protocol: Some("WebSocket".to_string()),
469        status_code: None,
470        status_text: None,
471        headers: headers.clone(),
472        content_length: Some(body.len()),
473        has_body: true,
474        is_chunked: false,
475        body_hex: Some(hex::encode(body.as_bytes())),
476        total_size: headers_size(headers) + body.len(),
477        original_source: "ssl.websocket".to_string(),
478        raw_data: include_raw_data.then(|| body.clone()),
479        body: Some(body),
480    }
481    .to_event(original_event)
482}
483
484impl HTTP2State {
485    fn handle_event(
486        &mut self,
487        original_event: &Event,
488        bytes: &[u8],
489        include_raw_data: bool,
490    ) -> Option<Vec<Event>> {
491        let tid = original_event
492            .data
493            .get("tid")
494            .and_then(|v| v.as_u64())
495            .unwrap_or(0);
496        let direction = direction_from_function(
497            original_event
498                .data
499                .get("function")
500                .and_then(|v| v.as_str())
501                .unwrap_or(""),
502        )?;
503        let frames = parse_http2_frames(bytes)?;
504        let mut events = Vec::new();
505
506        for frame in frames {
507            let key = (tid, frame.stream_id);
508            match frame.frame_type {
509                0x0 => {
510                    if frame.stream_id == 0 {
511                        continue;
512                    }
513                    let payload = data_payload(frame.flags, frame.payload);
514                    let state = self.streams.entry(key).or_default();
515                    match direction {
516                        HTTP2Direction::Request => {
517                            extend_capped(&mut state.request_body, payload, MAX_HTTP_BODY_BYTES);
518                            if frame.flags & 0x1 != 0 && !state.request_emitted {
519                                events.push(create_http2_request_event(
520                                    tid,
521                                    frame.stream_id,
522                                    state,
523                                    original_event,
524                                    include_raw_data,
525                                ));
526                                state.request_emitted = true;
527                            }
528                        }
529                        HTTP2Direction::Response => {
530                            extend_capped(&mut state.response_body, payload, MAX_HTTP_BODY_BYTES);
531                            if (frame.flags & 0x1 != 0
532                                || looks_like_complete_json(&state.response_body))
533                                && !state.response_emitted
534                            {
535                                events.push(create_http2_response_event(
536                                    tid,
537                                    frame.stream_id,
538                                    state,
539                                    original_event,
540                                    include_raw_data,
541                                ));
542                                state.response_emitted = true;
543                            }
544                        }
545                    }
546                }
547                0x1 => {
548                    if frame.stream_id == 0 {
549                        continue;
550                    }
551                    let fragment = headers_payload(frame.flags, frame.payload);
552                    if frame.flags & 0x4 != 0 {
553                        if let Some(headers) = self.decode_headers(direction, fragment) {
554                            let state = self.streams.entry(key).or_default();
555                            apply_headers(state, direction, headers);
556                            if frame.flags & 0x1 != 0 {
557                                match direction {
558                                    HTTP2Direction::Request if !state.request_emitted => {
559                                        events.push(create_http2_request_event(
560                                            tid,
561                                            frame.stream_id,
562                                            state,
563                                            original_event,
564                                            include_raw_data,
565                                        ));
566                                        state.request_emitted = true;
567                                    }
568                                    HTTP2Direction::Response if !state.response_emitted => {
569                                        events.push(create_http2_response_event(
570                                            tid,
571                                            frame.stream_id,
572                                            state,
573                                            original_event,
574                                            include_raw_data,
575                                        ));
576                                        state.response_emitted = true;
577                                    }
578                                    _ => {}
579                                }
580                            }
581                        }
582                    } else if fragment.len() <= MAX_HTTP2_HEADER_BLOCK_BYTES {
583                        self.pending_headers.insert(
584                            key,
585                            PendingHTTP2Headers {
586                                direction,
587                                block: fragment.to_vec(),
588                            },
589                        );
590                        evict_over_capacity(&mut self.pending_headers, MAX_HTTP2_PENDING_HEADERS);
591                    }
592                }
593                0x9 => {
594                    if frame.stream_id == 0 {
595                        continue;
596                    }
597                    let Some(mut pending) = self.pending_headers.remove(&key) else {
598                        continue;
599                    };
600                    pending.block.extend_from_slice(frame.payload);
601                    if pending.block.len() > MAX_HTTP2_HEADER_BLOCK_BYTES {
602                        continue;
603                    }
604                    if frame.flags & 0x4 != 0 {
605                        if let Some(headers) =
606                            self.decode_headers(pending.direction, &pending.block)
607                        {
608                            let state = self.streams.entry(key).or_default();
609                            apply_headers(state, pending.direction, headers);
610                        }
611                    } else {
612                        self.pending_headers.insert(key, pending);
613                    }
614                }
615                _ => {}
616            }
617
618            if self
619                .streams
620                .get(&key)
621                .map(|s| s.request_emitted && s.response_emitted)
622                .unwrap_or(false)
623            {
624                self.streams.remove(&key);
625            }
626            evict_over_capacity(&mut self.streams, MAX_HTTP2_STREAMS);
627        }
628
629        Some(if events.is_empty() {
630            Vec::new()
631        } else {
632            events
633        })
634    }
635
636    fn decode_headers(
637        &mut self,
638        direction: HTTP2Direction,
639        block: &[u8],
640    ) -> Option<HashMap<String, String>> {
641        let decoder = match direction {
642            HTTP2Direction::Request => &mut self.request_decoder,
643            HTTP2Direction::Response => &mut self.response_decoder,
644        };
645        let decoded = decoder.decode(block).ok()?;
646        let mut headers = HashMap::new();
647        for (name, value) in decoded {
648            let name = String::from_utf8_lossy(&name).to_ascii_lowercase();
649            let value = String::from_utf8_lossy(&value).to_string();
650            headers.insert(name, value);
651        }
652        if let Some(authority) = headers.get(":authority").cloned() {
653            headers.entry("host".to_string()).or_insert(authority);
654        }
655        Some(headers)
656    }
657}
658
659fn apply_headers(
660    state: &mut HTTP2StreamState,
661    direction: HTTP2Direction,
662    headers: HashMap<String, String>,
663) {
664    match direction {
665        HTTP2Direction::Request => state.request_headers.extend(headers),
666        HTTP2Direction::Response => state.response_headers.extend(headers),
667    }
668}
669
670fn direction_from_function(function: &str) -> Option<HTTP2Direction> {
671    let upper = function.to_ascii_uppercase();
672    if upper.contains("READ") || upper.contains("RECV") {
673        Some(HTTP2Direction::Response)
674    } else if upper.contains("WRITE") || upper.contains("SEND") {
675        Some(HTTP2Direction::Request)
676    } else {
677        None
678    }
679}
680
681fn parse_http2_frames(mut bytes: &[u8]) -> Option<Vec<HTTP2Frame<'_>>> {
682    const PREFACE: &[u8] = b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n";
683    if bytes.starts_with(PREFACE) {
684        bytes = &bytes[PREFACE.len()..];
685    }
686    if bytes.len() < 9 {
687        return None;
688    }
689
690    let mut frames = Vec::new();
691    let mut offset = 0usize;
692    while offset + 9 <= bytes.len() {
693        let length = ((bytes[offset] as usize) << 16)
694            | ((bytes[offset + 1] as usize) << 8)
695            | bytes[offset + 2] as usize;
696        let frame_type = bytes[offset + 3];
697        let flags = bytes[offset + 4];
698        let stream_id = ((bytes[offset + 5] as u32 & 0x7f) << 24)
699            | ((bytes[offset + 6] as u32) << 16)
700            | ((bytes[offset + 7] as u32) << 8)
701            | bytes[offset + 8] as u32;
702        offset += 9;
703        if length > bytes.len().saturating_sub(offset) {
704            return None;
705        }
706        let payload = &bytes[offset..offset + length];
707        offset += length;
708        // Skip unknown frame types per HTTP/2 spec (only process 0x0..=0x9)
709        if frame_type > 0x9 {
710            continue;
711        }
712        frames.push(HTTP2Frame {
713            frame_type,
714            flags,
715            stream_id,
716            payload,
717        });
718    }
719
720    if frames.is_empty() || offset != bytes.len() {
721        None
722    } else {
723        Some(frames)
724    }
725}
726
727fn headers_payload(flags: u8, payload: &[u8]) -> &[u8] {
728    let mut start = 0usize;
729    let mut end = payload.len();
730    if flags & 0x8 != 0 {
731        let Some(pad_len) = payload.first().copied() else {
732            return &[];
733        };
734        start += 1;
735        end = end.saturating_sub(pad_len as usize);
736    }
737    if flags & 0x20 != 0 {
738        start = start.saturating_add(5);
739    }
740    if start > end || end > payload.len() {
741        &[]
742    } else {
743        &payload[start..end]
744    }
745}
746
747fn data_payload(flags: u8, payload: &[u8]) -> &[u8] {
748    if flags & 0x8 == 0 {
749        return payload;
750    }
751    let Some(pad_len) = payload.first().copied() else {
752        return &[];
753    };
754    let start = 1usize;
755    let end = payload.len().saturating_sub(pad_len as usize);
756    if start > end || end > payload.len() {
757        &[]
758    } else {
759        &payload[start..end]
760    }
761}
762
763fn looks_like_complete_json(bytes: &[u8]) -> bool {
764    let text = String::from_utf8_lossy(bytes);
765    text.contains("usageMetadata") && serde_json::from_str::<serde_json::Value>(&text).is_ok()
766}
767
768fn create_http2_request_event(
769    tid: u64,
770    stream_id: u32,
771    state: &HTTP2StreamState,
772    original_event: &Event,
773    include_raw_data: bool,
774) -> Event {
775    let method = state.request_headers.get(":method").cloned();
776    let path = state.request_headers.get(":path").cloned();
777    let first_line = format!(
778        "{} {} HTTP/2",
779        method.as_deref().unwrap_or("HTTP"),
780        path.as_deref().unwrap_or("/")
781    );
782    let body = body_string(&state.request_body);
783    let body_hex = (!state.request_body.is_empty()).then(|| hex::encode(&state.request_body));
784    let total_size = headers_size(&state.request_headers) + state.request_body.len();
785    HTTPEvent {
786        tid: synthetic_http2_tid(tid, stream_id),
787        message_type: "request".to_string(),
788        first_line,
789        method,
790        path,
791        protocol: Some("HTTP/2".to_string()),
792        status_code: None,
793        status_text: None,
794        headers: state.request_headers.clone(),
795        content_length: body.as_ref().map(String::len),
796        has_body: body.is_some(),
797        is_chunked: false,
798        body,
799        body_hex,
800        total_size,
801        original_source: "ssl.http2".to_string(),
802        raw_data: include_raw_data
803            .then(|| String::from_utf8_lossy(&state.request_body).to_string()),
804    }
805    .to_event(original_event)
806}
807
808fn create_http2_response_event(
809    tid: u64,
810    stream_id: u32,
811    state: &HTTP2StreamState,
812    original_event: &Event,
813    include_raw_data: bool,
814) -> Event {
815    let status_code = state
816        .response_headers
817        .get(":status")
818        .and_then(|s| s.parse::<u16>().ok())
819        .or(Some(200));
820    let first_line = format!("HTTP/2 {}", status_code.unwrap_or(200));
821    let body = body_string(&state.response_body);
822    let body_hex = (!state.response_body.is_empty()).then(|| hex::encode(&state.response_body));
823    let total_size = headers_size(&state.response_headers) + state.response_body.len();
824    HTTPEvent {
825        tid: synthetic_http2_tid(tid, stream_id),
826        message_type: "response".to_string(),
827        first_line,
828        method: None,
829        path: None,
830        protocol: Some("HTTP/2".to_string()),
831        status_code,
832        status_text: None,
833        headers: state.response_headers.clone(),
834        content_length: body.as_ref().map(String::len),
835        has_body: body.is_some(),
836        is_chunked: false,
837        body,
838        body_hex,
839        total_size,
840        original_source: "ssl.http2".to_string(),
841        raw_data: include_raw_data
842            .then(|| String::from_utf8_lossy(&state.response_body).to_string()),
843    }
844    .to_event(original_event)
845}
846
847fn synthetic_http2_tid(tid: u64, stream_id: u32) -> u64 {
848    tid.saturating_mul(1_000_000)
849        .saturating_add(stream_id as u64)
850}
851
852fn body_string(body: &[u8]) -> Option<String> {
853    if body.is_empty() {
854        None
855    } else {
856        Some(String::from_utf8_lossy(body).to_string())
857    }
858}
859
860fn headers_size(headers: &HashMap<String, String>) -> usize {
861    headers.iter().map(|(k, v)| k.len() + v.len()).sum()
862}
863
864fn ssl_json_string_to_bytes(data: &str) -> Vec<u8> {
865    let mut bytes = Vec::with_capacity(data.len());
866    for ch in data.chars() {
867        let code = ch as u32;
868        if code <= 0xff {
869            bytes.push(code as u8);
870        } else {
871            let mut buf = [0u8; 4];
872            bytes.extend_from_slice(ch.encode_utf8(&mut buf).as_bytes());
873        }
874    }
875    bytes
876}
877
878#[async_trait]
879impl Analyzer for HTTPParser {
880    async fn process(&mut self, stream: EventStream) -> Result<EventStream, AnalyzerError> {
881        let include_raw_data = self.include_raw_data;
882        let mut http2 = std::mem::take(&mut self.http2);
883        let mut websocket = std::mem::take(&mut self.websocket);
884
885        let processed_stream = stream.flat_map(move |event| {
886            let events = if event.source == "ssl" {
887                Self::handle_ssl_event(&mut http2, &mut websocket, event, include_raw_data)
888            } else {
889                vec![event]
890            };
891            stream::iter(events)
892        });
893
894        Ok(Box::pin(processed_stream))
895    }
896}
897
898fn extend_capped(buffer: &mut Vec<u8>, data: &[u8], max: usize) {
899    buffer.extend_from_slice(data);
900    let overflow = buffer.len().saturating_sub(max);
901    if overflow > 0 {
902        buffer.drain(0..overflow);
903    }
904}
905
906fn evict_over_capacity<T>(map: &mut HashMap<(u64, u32), T>, max: usize) {
907    while map.len() > max {
908        let Some(key) = map.keys().next().copied() else {
909            break;
910        };
911        map.remove(&key);
912    }
913}
914
915#[cfg(test)]
916mod tests {
917    use super::*;
918    use crate::analyzers::{HTTPDecompressor, SSEProcessor};
919    use crate::view::MaterializedView;
920    use flate2::write::GzEncoder;
921    use flate2::{Compress, Compression, FlushCompress};
922    use futures::StreamExt;
923    use hpack::Encoder as HpackEncoder;
924    use serde_json::json;
925    use std::io::Write;
926
927    fn ssl_event(timestamp: u64, function: &str, bytes: Vec<u8>) -> Event {
928        Event::new_with_timestamp(
929            timestamp,
930            "ssl".to_string(),
931            4242,
932            "node".to_string(),
933            json!({
934                "tid": 7,
935                "function": function,
936                "data": bytes_to_ssl_json_string(&bytes),
937                "data_hex": hex::encode(&bytes),
938            }),
939        )
940    }
941
942    fn bytes_to_ssl_json_string(bytes: &[u8]) -> String {
943        bytes.iter().map(|b| char::from(*b)).collect()
944    }
945
946    fn frame(frame_type: u8, flags: u8, stream_id: u32, payload: &[u8]) -> Vec<u8> {
947        let len = payload.len();
948        let mut out = vec![
949            ((len >> 16) & 0xff) as u8,
950            ((len >> 8) & 0xff) as u8,
951            (len & 0xff) as u8,
952            frame_type,
953            flags,
954            ((stream_id >> 24) & 0x7f) as u8,
955            ((stream_id >> 16) & 0xff) as u8,
956            ((stream_id >> 8) & 0xff) as u8,
957            (stream_id & 0xff) as u8,
958        ];
959        out.extend_from_slice(payload);
960        out
961    }
962
963    fn compressed_websocket_frame(compressor: &mut Compress, payload: &[u8]) -> Vec<u8> {
964        let mut compressed = Vec::with_capacity(payload.len() * 2 + 64);
965        compressor
966            .compress_vec(payload, &mut compressed, FlushCompress::Sync)
967            .unwrap();
968        assert!(compressed.ends_with(&[0, 0, 0xff, 0xff]));
969        compressed.truncate(compressed.len() - 4);
970
971        let mut frame = vec![0xc1];
972        if compressed.len() <= 125 {
973            frame.push(0x80 | compressed.len() as u8);
974        } else {
975            frame.push(0x80 | 126);
976            frame.extend_from_slice(&(compressed.len() as u16).to_be_bytes());
977        }
978        let mask = [0x12, 0x34, 0x56, 0x78];
979        frame.extend_from_slice(&mask);
980        frame.extend(
981            compressed
982                .iter()
983                .enumerate()
984                .map(|(i, byte)| byte ^ mask[i % 4]),
985        );
986        frame
987    }
988
989    #[tokio::test]
990    async fn parses_compressed_websocket_responses_with_context_takeover() {
991        let handshake = b"GET /backend-api/codex/responses HTTP/1.1\r\n\
992Host: chatgpt.com\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\
993sec-websocket-extensions: permessage-deflate\r\n\r\n"
994            .to_vec();
995        let shared = "shared context ".repeat(400);
996        let first = json!({
997            "type": "response.create",
998            "model": "gpt-test",
999            "input": [{"role": "developer", "content": shared}],
1000        })
1001        .to_string();
1002        let prompt = "agentsight websocket exact prompt 7f31";
1003        let second = json!({
1004            "type": "response.create",
1005            "model": "gpt-test",
1006            "input": [{"role": "user", "content": format!("{shared}{prompt}")}],
1007        })
1008        .to_string();
1009        let mut compressor = Compress::new(Compression::fast(), false);
1010        let input: EventStream = Box::pin(stream::iter(vec![
1011            ssl_event(1, "WRITE/SEND", handshake),
1012            ssl_event(
1013                2,
1014                "WRITE/SEND",
1015                compressed_websocket_frame(&mut compressor, first.as_bytes()),
1016            ),
1017            ssl_event(
1018                3,
1019                "WRITE/SEND",
1020                compressed_websocket_frame(&mut compressor, second.as_bytes()),
1021            ),
1022        ]));
1023        let mut parser = HTTPParser::new().disable_raw_data();
1024        let output: Vec<Event> = parser.process(input).await.unwrap().collect().await;
1025
1026        assert_eq!(output.len(), 3);
1027        assert_eq!(output[2].data["path"], "/backend-api/codex/responses");
1028        assert!(output[2].data["body"].as_str().unwrap().contains(prompt));
1029        let mut view = MaterializedView::new();
1030        for event in output {
1031            view.ingest_event(&event).unwrap();
1032        }
1033        let calls = view.llm_call_rows(10);
1034        assert_eq!(calls.len(), 2);
1035        assert_eq!(calls[0].call_kind.as_deref(), Some("responses"));
1036        assert!(
1037            calls
1038                .iter()
1039                .any(|call| call.request.to_string().contains(prompt))
1040        );
1041    }
1042
1043    #[tokio::test]
1044    async fn parses_http2_gemini_usage_into_http_events() {
1045        let mut request_encoder = HpackEncoder::new();
1046        let mut response_encoder = HpackEncoder::new();
1047        let request_headers = [
1048            (&b":method"[..], &b"POST"[..]),
1049            (&b":scheme"[..], &b"https"[..]),
1050            (&b":authority"[..], &b"cloudcode-pa.googleapis.com"[..]),
1051            (&b":path"[..], &b"/v1internal:generateContent"[..]),
1052        ];
1053        let response_headers = [
1054            (&b":status"[..], &b"200"[..]),
1055            (&b"content-type"[..], &b"application/json"[..]),
1056        ];
1057        let request_body = br#"{"model":"gemini-2.5-pro","request":{"contents":[]}}"#;
1058        let response_body = br#"{"usageMetadata":{"promptTokenCount":11,"candidatesTokenCount":4,"totalTokenCount":15}}"#;
1059
1060        let mut request_bytes = b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n".to_vec();
1061        request_bytes.extend(frame(0x1, 0x4, 1, &request_encoder.encode(request_headers)));
1062        request_bytes.extend(frame(0x0, 0x1, 1, request_body));
1063
1064        let mut response_bytes = Vec::new();
1065        response_bytes.extend(frame(
1066            0x1,
1067            0x4,
1068            1,
1069            &response_encoder.encode(response_headers),
1070        ));
1071        response_bytes.extend(frame(0x0, 0x1, 1, response_body));
1072
1073        let input: EventStream = Box::pin(stream::iter(vec![
1074            ssl_event(1, "WRITE/SEND", request_bytes),
1075            ssl_event(2, "READ/RECV", response_bytes),
1076        ]));
1077        let mut parser = HTTPParser::new().disable_raw_data();
1078        let output: Vec<Event> = parser.process(input).await.unwrap().collect().await;
1079
1080        assert_eq!(output.len(), 2);
1081        assert_eq!(output[0].source, "http_parser");
1082        assert_eq!(output[0].data["message_type"], "request");
1083        assert_eq!(output[0].data["path"], "/v1internal:generateContent");
1084        assert_eq!(
1085            output[0].data["headers"]["host"],
1086            "cloudcode-pa.googleapis.com"
1087        );
1088        assert_eq!(output[1].source, "http_parser");
1089        assert_eq!(output[1].data["message_type"], "response");
1090        assert_eq!(output[1].data["status_code"], 200);
1091        assert!(
1092            output[1].data["body"]
1093                .as_str()
1094                .unwrap()
1095                .contains("usageMetadata")
1096        );
1097
1098        let mut view = MaterializedView::new();
1099        for event in output {
1100            view.ingest_event(&event).unwrap();
1101        }
1102        let total = view
1103            .export_snapshot(crate::model::SnapshotOptions { audit_limit: 0 })
1104            .token_summary
1105            .into_iter()
1106            .map(|row| row.total_tokens)
1107            .sum::<i64>();
1108        assert_eq!(total, 15);
1109    }
1110
1111    #[tokio::test]
1112    async fn http2_gzip_sse_capture_pipeline_reaches_materialized_view() {
1113        let mut request_encoder = HpackEncoder::new();
1114        let mut response_encoder = HpackEncoder::new();
1115        let request_headers = [
1116            (&b":method"[..], &b"POST"[..]),
1117            (&b":scheme"[..], &b"https"[..]),
1118            (&b":authority"[..], &b"api.openai.com"[..]),
1119            (&b":path"[..], &b"/v1/chat/completions"[..]),
1120        ];
1121        let response_headers = [
1122            (&b":status"[..], &b"200"[..]),
1123            (&b"content-type"[..], &b"text/event-stream"[..]),
1124            (&b"content-encoding"[..], &b"gzip"[..]),
1125        ];
1126        let request_body = br#"{"model":"gpt-test","metadata":{"session_id":"sess-h2"}}"#;
1127        let mut gzip = GzEncoder::new(Vec::new(), Compression::default());
1128        gzip.write_all(
1129            b"data: {\"choices\":[{\"finish_reason\":\"stop\"}],\"usage\":{\"prompt_tokens\":2,\"completion_tokens\":5}}\n\ndata: [DONE]\n\n",
1130        )
1131        .unwrap();
1132        let response_body = gzip.finish().unwrap();
1133
1134        let mut request_bytes = b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n".to_vec();
1135        request_bytes.extend(frame(0x1, 0x4, 1, &request_encoder.encode(request_headers)));
1136        request_bytes.extend(frame(0x0, 0x1, 1, request_body));
1137
1138        let mut response_bytes = Vec::new();
1139        response_bytes.extend(frame(
1140            0x1,
1141            0x4,
1142            1,
1143            &response_encoder.encode(response_headers),
1144        ));
1145        response_bytes.extend(frame(0x0, 0x1, 1, &response_body));
1146
1147        let input: EventStream = Box::pin(stream::iter(vec![
1148            ssl_event(1, "WRITE/SEND", request_bytes),
1149            ssl_event(2, "READ/RECV", response_bytes),
1150        ]));
1151        let mut parser = HTTPParser::new().disable_raw_data();
1152        let parsed = parser.process(input).await.unwrap();
1153        let mut decompressor = HTTPDecompressor::new();
1154        let decompressed = decompressor.process(parsed).await.unwrap();
1155        let mut sse = SSEProcessor::new();
1156        let output: Vec<Event> = sse.process(decompressed).await.unwrap().collect().await;
1157
1158        assert_eq!(output.len(), 2);
1159        assert_eq!(output[0].data["message_type"], "request");
1160        assert_eq!(output[1].source, "sse_processor");
1161        assert_eq!(output[1].data["status_code"], 200);
1162
1163        let mut view = MaterializedView::new();
1164        for event in output {
1165            view.ingest_event(&event).unwrap();
1166        }
1167        let snapshot = view.export_snapshot(crate::model::SnapshotOptions { audit_limit: 0 });
1168        assert_eq!(snapshot.summary.llm_calls, 1);
1169        assert_eq!(snapshot.summary.token_usage_rows, 1);
1170        assert_eq!(snapshot.summary.input_tokens, 2);
1171        assert_eq!(snapshot.summary.output_tokens, 5);
1172        assert_eq!(snapshot.summary.total_tokens, 7);
1173        let calls = view.llm_call_rows(10);
1174        assert_eq!(calls[0].status, "complete");
1175        assert_eq!(calls[0].session_id.as_deref(), Some("sess-h2"));
1176        assert_eq!(calls[0].call_kind.as_deref(), Some("chat"));
1177        assert_eq!(calls[0].finish_reason.as_deref(), Some("stop"));
1178    }
1179
1180    #[test]
1181    fn rejects_non_http2_frames() {
1182        assert!(parse_http2_frames(b"GET / HTTP/1.1\r\n\r\n").is_none());
1183    }
1184}