1use 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
19pub struct HTTPParser {
21 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#[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 pub method: Option<String>,
100 pub path: Option<String>,
101 pub protocol: Option<String>,
102 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 pub fn new() -> Self {
116 HTTPParser {
117 include_raw_data: true,
118 http2: HTTP2State::default(),
119 websocket: WebSocketState::default(),
120 }
121 }
122
123 pub fn disable_raw_data(mut self) -> Self {
125 self.include_raw_data = false;
126 self
127 }
128
129 pub fn is_http_data(data: &str) -> bool {
131 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 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 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 if first_line.starts_with("HTTP/") {
174 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 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 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 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 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 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 let total_size = parsed_message.first_line.len() +
278 parsed_message.headers.iter().map(|(k, v)| k.len() + v.len() + 4).sum::<usize>() + parsed_message.body.as_ref().map(|b| b.len()).unwrap_or(0) +
280 4; 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 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 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 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 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}