Skip to main content

agentsight_capture/analyzers/
sse_processor.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 futures::stream::StreamExt;
9use serde_json::{Value, json};
10use std::collections::{BTreeMap, HashMap};
11use std::sync::{Arc, Mutex};
12
13use super::protocol_events::SSEProcessorEvent;
14
15const MAX_BUFFERS: usize = 1024;
16
17pub struct SSEProcessor {
18    sse_buffers: Arc<Mutex<HashMap<String, SSEAccumulator>>>,
19    timeout_ms: u64,
20    max_buffers: usize,
21}
22
23impl Default for SSEProcessor {
24    fn default() -> Self {
25        Self::new_with_timeout(30_000)
26    }
27}
28
29struct SSEAccumulator {
30    message_id: Option<String>,
31    accumulated_text: String,
32    accumulated_json: String,
33    openai_reasoning: String,
34    openai_tool_calls: BTreeMap<u64, OpenAIToolCallAccumulator>,
35    events: Vec<SSEEvent>,
36    is_complete: bool,
37    last_update: u64,
38    has_message_start: bool,
39    start_time: u64,
40    end_time: u64,
41}
42
43#[derive(Clone, Debug, Default)]
44struct OpenAIToolCallAccumulator {
45    id: Option<String>,
46    type_name: Option<String>,
47    function_name: Option<String>,
48    function_arguments: String,
49}
50
51#[derive(Clone, Debug)]
52pub struct SSEEvent {
53    pub event: Option<String>,
54    pub data: Option<String>,
55    pub id: Option<String>,
56    pub parsed_data: Option<Value>,
57    pub raw_data: Option<String>,
58}
59
60impl SSEProcessor {
61    #[cfg(test)]
62    pub fn new() -> Self {
63        Self::default()
64    }
65
66    pub fn new_with_timeout(timeout_ms: u64) -> Self {
67        SSEProcessor {
68            sse_buffers: Arc::new(Mutex::new(HashMap::new())),
69            timeout_ms,
70            max_buffers: MAX_BUFFERS,
71        }
72    }
73
74    pub fn is_sse_data(data: &str) -> bool {
75        let has_sse_patterns = data.contains("event:") && data.contains("data:");
76        let has_sse_content_type = data.contains("text/event-stream");
77        let has_chunked_sse = data.contains("Transfer-Encoding: chunked")
78            && (data.contains("event:") || data.contains("data:"));
79        let has_sse_data_only =
80            data.contains("data:") && (data.contains("\r\n\r\n") || data.contains("\n\n"));
81        has_sse_patterns || has_sse_content_type || has_chunked_sse || has_sse_data_only
82    }
83
84    pub fn parse_sse_events_from_chunk(chunk_content: &str) -> Vec<SSEEvent> {
85        let mut events = Vec::new();
86        let normalized = chunk_content.replace("\r\n", "\n");
87        let event_blocks: Vec<&str> = normalized.split("\n\n").collect();
88
89        for block in event_blocks {
90            if block.trim().is_empty() {
91                continue;
92            }
93
94            let mut event = SSEEvent {
95                event: None,
96                data: None,
97                id: None,
98                parsed_data: None,
99                raw_data: None,
100            };
101            let mut data_lines = Vec::new();
102
103            for line in block.split('\n') {
104                let line = line.trim();
105                if let Some(rest) = line.strip_prefix("event:") {
106                    event.event = Some(rest.trim().to_string());
107                } else if let Some(rest) = line.strip_prefix("data:") {
108                    data_lines.push(rest.trim());
109                } else if let Some(rest) = line.strip_prefix("id:") {
110                    event.id = Some(rest.trim().to_string());
111                }
112            }
113
114            if !data_lines.is_empty() {
115                let combined_data = data_lines.join("\n");
116                event.data = Some(combined_data.clone());
117                match serde_json::from_str::<Value>(&combined_data) {
118                    Ok(parsed_json) => {
119                        event.parsed_data = Some(parsed_json);
120                    }
121                    Err(_) => {
122                        event.raw_data = Some(combined_data);
123                    }
124                }
125            }
126
127            if event.event.is_some() || event.data.is_some() {
128                events.push(event);
129            }
130        }
131
132        events
133    }
134
135    pub fn parse_sse_events(data: &str) -> Vec<SSEEvent> {
136        let clean_data = Self::clean_chunked_content(data);
137        let sse_data = if clean_data.trim().is_empty() {
138            data
139        } else {
140            clean_data.as_str()
141        };
142        Self::parse_sse_events_from_chunk(sse_data)
143    }
144
145    fn sse_payload(event: &Event) -> Option<(&str, bool)> {
146        if event.source == "ssl" {
147            return event
148                .data
149                .get("data")
150                .and_then(|v| v.as_str())
151                .map(|data| (data, true));
152        }
153
154        if event.source != "http_parser"
155            || event.data.get("message_type").and_then(|v| v.as_str()) != Some("response")
156        {
157            return None;
158        }
159
160        let body = event.data.get("body").and_then(|v| v.as_str())?;
161        (Self::http_content_type_is_sse(&event.data) || Self::is_sse_data(body))
162            .then_some((body, false))
163    }
164
165    fn http_content_type_is_sse(data: &Value) -> bool {
166        let Some(headers) = data.get("headers").and_then(|v| v.as_object()) else {
167            return false;
168        };
169        headers.iter().any(|(key, value)| {
170            key.eq_ignore_ascii_case("content-type")
171                && value
172                    .as_str()
173                    .is_some_and(|v| v.to_ascii_lowercase().contains("text/event-stream"))
174        })
175    }
176
177    fn parse_usage_metadata_fragment(data: &str) -> Option<SSEEvent> {
178        let usage = extract_json_object_after_key(data, "\"usageMetadata\"")?;
179        let usage_json: Value = serde_json::from_str(usage).ok()?;
180        let has_tokens = usage_json.get("promptTokenCount").is_some()
181            || usage_json.get("candidatesTokenCount").is_some()
182            || usage_json.get("totalTokenCount").is_some();
183        if !has_tokens {
184            return None;
185        }
186
187        let mut parsed = serde_json::Map::new();
188        parsed.insert("usageMetadata".to_string(), usage_json);
189        if let Some(model) = extract_json_string_field(data, "modelVersion")
190            .or_else(|| extract_json_string_field(data, "model"))
191        {
192            parsed.insert("modelVersion".to_string(), Value::String(model));
193        }
194
195        Some(SSEEvent {
196            event: Some("message_stop".to_string()),
197            data: None,
198            id: None,
199            parsed_data: Some(Value::Object(parsed)),
200            raw_data: None,
201        })
202    }
203
204    pub fn clean_chunked_content(content: &str) -> String {
205        let mut content_parts = Vec::new();
206        let lines: Vec<&str> = content.split("\r\n").collect();
207
208        let mut i = 0;
209        while i < lines.len() {
210            let line = lines[i].trim();
211            if !line.is_empty() && line.chars().all(|c| c.is_ascii_hexdigit()) {
212                let chunk_size = u32::from_str_radix(line, 16).unwrap_or(0);
213                if chunk_size == 0 {
214                    break;
215                }
216                i += 1;
217                if i < lines.len() {
218                    content_parts.push(lines[i]);
219                }
220            }
221            i += 1;
222        }
223
224        content_parts.join("\n")
225    }
226
227    fn generate_connection_id(event: &Event, sse_events: &[SSEEvent]) -> String {
228        let pid = event.data.get("pid").and_then(|v| v.as_u64()).unwrap_or(0);
229        let tid = event.data.get("tid").and_then(|v| v.as_u64()).unwrap_or(0);
230
231        if let Some(message_id) = Self::extract_message_id(sse_events) {
232            return format!("{}:{}:{}", pid, tid, message_id);
233        }
234
235        let timestamp = event.timestamp;
236        let window = timestamp / 600_000_000_000;
237        format!("{}:{}:{}", pid, tid, window)
238    }
239
240    fn extract_message_id(events: &[SSEEvent]) -> Option<String> {
241        for event in events {
242            if let Some(event_type) = &event.event
243                && event_type == "message_start"
244                && let Some(parsed_data) = &event.parsed_data
245                && let Some(message) = parsed_data.get("message")
246                && let Some(id) = message.get("id")
247                && let Some(id_str) = id.as_str()
248            {
249                return Some(id_str.to_string());
250            }
251        }
252        for event in events {
253            if let Some(parsed_data) = &event.parsed_data
254                && let Some(id) = parsed_data.get("id").and_then(|id| id.as_str())
255            {
256                return Some(id.to_string());
257            }
258            if let Some(parsed_data) = &event.parsed_data {
259                if let Some(id) = parsed_data.get("response_id").and_then(|id| id.as_str()) {
260                    return Some(id.to_string());
261                }
262                if let Some(id) = parsed_data
263                    .get("response")
264                    .and_then(|response| response.get("id"))
265                    .and_then(|id| id.as_str())
266                {
267                    return Some(id.to_string());
268                }
269            }
270        }
271        None
272    }
273
274    fn is_sse_complete(accumulator: &SSEAccumulator) -> bool {
275        for event in &accumulator.events {
276            if Self::sse_event_completes_stream(event) {
277                return true;
278            }
279            if let Some(event_type) = &event.event {
280                match event_type.as_str() {
281                    "message_stop" => return true,
282                    "error" => return true,
283                    _ => {}
284                }
285            }
286        }
287        accumulator.accumulated_text.len() > 50000 || accumulator.accumulated_json.len() > 50000
288    }
289
290    fn has_meaningful_content(accumulator: &SSEAccumulator) -> bool {
291        if !accumulator.accumulated_text.is_empty() || !accumulator.accumulated_json.is_empty() {
292            return true;
293        }
294        if !accumulator.openai_reasoning.is_empty() || !accumulator.openai_tool_calls.is_empty() {
295            return true;
296        }
297
298        let mut has_content_deltas = false;
299        let mut has_message_start = false;
300        let mut metadata_only_count = 0;
301
302        for event in &accumulator.events {
303            if Self::sse_event_has_usage(event) {
304                return true;
305            }
306            if Self::sse_event_has_openai_delta(event) || Self::sse_event_has_terminal_finish(event)
307            {
308                return true;
309            }
310            if Self::sse_event_has_openai_response_delta(event)
311                || Self::sse_event_has_openai_response_terminal(event)
312            {
313                return true;
314            }
315            if let Some(event_type) = &event.event {
316                match event_type.as_str() {
317                    "content_block_delta" => has_content_deltas = true,
318                    "message_start" => has_message_start = true,
319                    "message_stop"
320                    | "message_delta"
321                    | "ping"
322                    | "content_block_stop"
323                    | "content_block_start" => {
324                        metadata_only_count += 1;
325                    }
326                    _ => {}
327                }
328            }
329        }
330
331        has_content_deltas
332            || (has_message_start
333                && accumulator.events.len() > 3
334                && metadata_only_count < accumulator.events.len())
335    }
336
337    fn sse_event_has_usage(event: &SSEEvent) -> bool {
338        event
339            .parsed_data
340            .as_ref()
341            .is_some_and(|data| Self::meaningful_usage(data).is_some())
342    }
343
344    fn sse_event_completes_stream(event: &SSEEvent) -> bool {
345        event.data.as_deref() == Some("[DONE]")
346            || event.parsed_data.as_ref().is_some_and(|data| {
347                Self::has_stream_completing_usage(data) || Self::has_openai_response_terminal(data)
348            })
349    }
350
351    fn has_stream_completing_usage(data: &Value) -> bool {
352        [
353            data.get("usageMetadata"),
354            data.get("usage"),
355            data.get("response")
356                .and_then(|response| response.get("usage")),
357        ]
358        .into_iter()
359        .flatten()
360        .any(Self::usage_has_meaningful_fields)
361    }
362
363    fn meaningful_usage(data: &Value) -> Option<&Value> {
364        [
365            data.get("usageMetadata"),
366            data.get("usage"),
367            data.get("message").and_then(|m| m.get("usage")),
368            data.get("response")
369                .and_then(|response| response.get("usage")),
370        ]
371        .into_iter()
372        .flatten()
373        .find(|usage| Self::usage_has_meaningful_fields(usage))
374    }
375
376    fn usage_has_meaningful_fields(usage: &Value) -> bool {
377        usage
378            .as_object()
379            .is_some_and(|fields| fields.values().any(|value| !value.is_null()))
380    }
381
382    fn sse_event_has_openai_delta(event: &SSEEvent) -> bool {
383        event
384            .parsed_data
385            .as_ref()
386            .and_then(|data| data.get("choices"))
387            .and_then(|choices| choices.as_array())
388            .is_some_and(|choices| {
389                choices.iter().any(|choice| {
390                    choice
391                        .get("delta")
392                        .and_then(|delta| delta.as_object())
393                        .is_some_and(|delta| {
394                            delta
395                                .get("content")
396                                .and_then(|v| v.as_str())
397                                .is_some_and(|v| !v.is_empty())
398                                || delta
399                                    .get("tool_calls")
400                                    .and_then(|v| v.as_array())
401                                    .is_some_and(|v| !v.is_empty())
402                                || delta.get("function_call").is_some()
403                                || Self::openai_reasoning_delta(delta).is_some()
404                        })
405                })
406            })
407    }
408
409    fn sse_event_has_openai_response_delta(event: &SSEEvent) -> bool {
410        event
411            .parsed_data
412            .as_ref()
413            .is_some_and(Self::has_openai_response_delta)
414    }
415
416    fn has_openai_response_delta(data: &Value) -> bool {
417        matches!(
418            Self::openai_response_event_type(data),
419            Some(
420                "response.output_text.delta"
421                    | "response.reasoning_text.delta"
422                    | "response.reasoning_summary_text.delta"
423                    | "response.function_call_arguments.delta"
424                    | "response.function_call_arguments.done"
425                    | "response.output_item.added"
426                    | "response.output_item.done"
427            )
428        )
429    }
430
431    fn sse_event_has_terminal_finish(event: &SSEEvent) -> bool {
432        event
433            .parsed_data
434            .as_ref()
435            .is_some_and(Self::has_terminal_finish_reason)
436    }
437
438    fn has_terminal_finish_reason(data: &Value) -> bool {
439        data.get("choices")
440            .and_then(|choices| choices.as_array())
441            .is_some_and(|choices| {
442                choices.iter().any(|choice| {
443                    choice
444                        .get("finish_reason")
445                        .is_some_and(|reason| !reason.is_null())
446                })
447            })
448    }
449
450    fn sse_event_has_openai_response_terminal(event: &SSEEvent) -> bool {
451        event
452            .parsed_data
453            .as_ref()
454            .is_some_and(Self::has_openai_response_terminal)
455    }
456
457    fn has_openai_response_terminal(data: &Value) -> bool {
458        matches!(
459            Self::openai_response_event_type(data),
460            Some(
461                "response.completed"
462                    | "response.failed"
463                    | "response.incomplete"
464                    | "response.cancelled"
465            )
466        )
467    }
468
469    fn openai_response_event_type(data: &Value) -> Option<&str> {
470        data.get("type")
471            .and_then(|value| value.as_str())
472            .filter(|value| value.starts_with("response."))
473    }
474
475    fn openai_reasoning_delta(delta: &serde_json::Map<String, Value>) -> Option<&str> {
476        [
477            "reasoning_content",
478            "reasoning",
479            "reasoning_text",
480            "thinking",
481        ]
482        .into_iter()
483        .find_map(|key| delta.get(key).and_then(|v| v.as_str()))
484        .filter(|value| !value.is_empty())
485    }
486
487    fn accumulate_content(accumulator: &mut SSEAccumulator, events: &[SSEEvent]) {
488        for event in events {
489            accumulator.events.push(event.clone());
490
491            if accumulator.message_id.is_none() {
492                accumulator.message_id = Self::extract_message_id(std::slice::from_ref(event));
493            }
494
495            if let Some(event_type) = &event.event {
496                match event_type.as_str() {
497                    "message_start" => {
498                        accumulator.has_message_start = true;
499                        if accumulator.message_id.is_none() {
500                            accumulator.message_id =
501                                Self::extract_message_id(std::slice::from_ref(event));
502                        }
503                    }
504                    "content_block_delta" => {
505                        if let Some(parsed_data) = &event.parsed_data
506                            && let Some(delta) = parsed_data.get("delta")
507                        {
508                            let delta_type = delta.get("type").and_then(|v| v.as_str());
509                            let text = if delta_type == Some("text_delta") {
510                                delta.get("text").and_then(|v| v.as_str())
511                            } else if delta_type == Some("thinking_delta") {
512                                delta.get("thinking").and_then(|v| v.as_str())
513                            } else {
514                                None
515                            };
516                            if let Some(t) = text {
517                                accumulator.accumulated_text.push_str(t);
518                            }
519                            if let Some(partial_json) =
520                                delta.get("partial_json").and_then(|v| v.as_str())
521                            {
522                                accumulator.accumulated_json.push_str(partial_json);
523                            }
524                        }
525                    }
526                    _ => {}
527                }
528            }
529            if let Some(parsed_data) = &event.parsed_data {
530                Self::accumulate_openai_content(accumulator, parsed_data);
531                Self::accumulate_openai_responses_content(accumulator, parsed_data);
532            }
533        }
534    }
535
536    fn accumulate_openai_content(accumulator: &mut SSEAccumulator, data: &Value) {
537        let Some(choices) = data.get("choices").and_then(|choices| choices.as_array()) else {
538            return;
539        };
540        for choice in choices {
541            let Some(delta) = choice.get("delta").and_then(|delta| delta.as_object()) else {
542                continue;
543            };
544            if let Some(content) = delta.get("content").and_then(|v| v.as_str()) {
545                accumulator.accumulated_text.push_str(content);
546            }
547            if let Some(reasoning) = Self::openai_reasoning_delta(delta) {
548                accumulator.openai_reasoning.push_str(reasoning);
549            }
550            if let Some(tool_calls) = delta.get("tool_calls").and_then(|v| v.as_array()) {
551                for (fallback_index, tool_call) in tool_calls.iter().enumerate() {
552                    Self::accumulate_openai_tool_call(
553                        accumulator,
554                        tool_call,
555                        fallback_index as u64,
556                    );
557                }
558            }
559            if let Some(function_call) = delta.get("function_call") {
560                Self::accumulate_openai_function_call(accumulator, function_call);
561            }
562        }
563    }
564
565    fn accumulate_openai_responses_content(accumulator: &mut SSEAccumulator, data: &Value) {
566        match Self::openai_response_event_type(data) {
567            Some("response.output_text.delta") => {
568                if let Some(delta) = data.get("delta").and_then(|value| value.as_str()) {
569                    accumulator.accumulated_text.push_str(delta);
570                }
571            }
572            Some("response.reasoning_text.delta" | "response.reasoning_summary_text.delta") => {
573                if let Some(delta) = data.get("delta").and_then(|value| value.as_str()) {
574                    accumulator.openai_reasoning.push_str(delta);
575                }
576            }
577            Some("response.function_call_arguments.delta") => {
578                Self::accumulate_openai_response_function_arguments(accumulator, data, false);
579            }
580            Some("response.function_call_arguments.done") => {
581                Self::accumulate_openai_response_function_arguments(accumulator, data, true);
582            }
583            Some("response.output_item.added" | "response.output_item.done") => {
584                Self::accumulate_openai_response_item(accumulator, data);
585            }
586            _ => {}
587        }
588    }
589
590    fn accumulate_openai_response_item(accumulator: &mut SSEAccumulator, data: &Value) {
591        let Some(item) = data.get("item").and_then(|value| value.as_object()) else {
592            return;
593        };
594        if item.get("type").and_then(|value| value.as_str()) != Some("function_call") {
595            return;
596        }
597        let index = data
598            .get("output_index")
599            .and_then(|value| value.as_u64())
600            .unwrap_or(accumulator.openai_tool_calls.len() as u64);
601        let entry = accumulator.openai_tool_calls.entry(index).or_default();
602        entry
603            .type_name
604            .get_or_insert_with(|| "function".to_string());
605        if let Some(id) = item.get("id").and_then(|value| value.as_str()) {
606            entry.id = Some(id.to_string());
607        } else if let Some(id) = item.get("call_id").and_then(|value| value.as_str()) {
608            entry.id = Some(id.to_string());
609        }
610        if let Some(name) = item.get("name").and_then(|value| value.as_str()) {
611            entry.function_name = Some(name.to_string());
612        }
613        if let Some(arguments) = item.get("arguments").and_then(|value| value.as_str()) {
614            entry.function_arguments = arguments.to_string();
615        }
616    }
617
618    fn accumulate_openai_response_function_arguments(
619        accumulator: &mut SSEAccumulator,
620        data: &Value,
621        replace: bool,
622    ) {
623        let index = data
624            .get("output_index")
625            .and_then(|value| value.as_u64())
626            .unwrap_or(0);
627        let entry = accumulator.openai_tool_calls.entry(index).or_default();
628        entry
629            .type_name
630            .get_or_insert_with(|| "function".to_string());
631        if let Some(id) = data.get("item_id").and_then(|value| value.as_str()) {
632            entry.id = Some(id.to_string());
633        }
634        let payload_key = if replace { "arguments" } else { "delta" };
635        if let Some(arguments) = data.get(payload_key).and_then(|value| value.as_str()) {
636            if replace {
637                entry.function_arguments = arguments.to_string();
638            } else {
639                entry.function_arguments.push_str(arguments);
640            }
641        }
642    }
643
644    fn accumulate_openai_tool_call(
645        accumulator: &mut SSEAccumulator,
646        tool_call: &Value,
647        fallback_index: u64,
648    ) {
649        let index = tool_call
650            .get("index")
651            .and_then(|v| v.as_u64())
652            .unwrap_or(fallback_index);
653        let entry = accumulator.openai_tool_calls.entry(index).or_default();
654        if let Some(id) = tool_call.get("id").and_then(|v| v.as_str()) {
655            entry.id = Some(id.to_string());
656        }
657        if let Some(type_name) = tool_call.get("type").and_then(|v| v.as_str()) {
658            entry.type_name = Some(type_name.to_string());
659        }
660        if let Some(function) = tool_call.get("function").and_then(|v| v.as_object()) {
661            if let Some(name) = function.get("name").and_then(|v| v.as_str()) {
662                entry.function_name = Some(name.to_string());
663            }
664            if let Some(arguments) = function.get("arguments").and_then(|v| v.as_str()) {
665                entry.function_arguments.push_str(arguments);
666            }
667        }
668    }
669
670    fn accumulate_openai_function_call(accumulator: &mut SSEAccumulator, function_call: &Value) {
671        let entry = accumulator.openai_tool_calls.entry(0).or_default();
672        entry
673            .type_name
674            .get_or_insert_with(|| "function".to_string());
675        if let Some(name) = function_call.get("name").and_then(|v| v.as_str()) {
676            entry.function_name = Some(name.to_string());
677        }
678        if let Some(arguments) = function_call.get("arguments").and_then(|v| v.as_str()) {
679            entry.function_arguments.push_str(arguments);
680        }
681    }
682
683    fn create_merged_event(
684        connection_id: String,
685        accumulator: &SSEAccumulator,
686        original_event: &Event,
687    ) -> Event {
688        let json_content = Self::merged_json_content(accumulator);
689
690        let text_content = accumulator.accumulated_text.clone();
691
692        let sse_events_json: Vec<Value> = accumulator
693            .events
694            .iter()
695            .map(|e| {
696                json!({
697                    "event": e.event,
698                    "data": e.data,
699                    "id": e.id,
700                    "parsed_data": e.parsed_data,
701                    "raw_data": e.raw_data
702                })
703            })
704            .collect();
705
706        let total_size = json_content.len() + text_content.len();
707
708        SSEProcessorEvent {
709            connection_id,
710            message_id: accumulator.message_id.clone(),
711            start_time: accumulator.start_time,
712            end_time: accumulator.end_time,
713            duration_ns: accumulator.end_time.saturating_sub(accumulator.start_time),
714            original_source: original_event.source.clone(),
715            host: Self::event_host(original_event),
716            method: original_event
717                .data
718                .get("method")
719                .and_then(|v| v.as_str())
720                .map(str::to_string),
721            path: original_event
722                .data
723                .get("path")
724                .and_then(|v| v.as_str())
725                .map(str::to_string),
726            status_code: original_event
727                .data
728                .get("status_code")
729                .and_then(|v| v.as_u64())
730                .map(|v| v as u16),
731            function: original_event
732                .data
733                .get("function")
734                .and_then(|v| v.as_str())
735                .unwrap_or("unknown")
736                .to_string(),
737            tid: original_event
738                .data
739                .get("tid")
740                .and_then(|v| v.as_u64())
741                .unwrap_or(0),
742            json_content,
743            text_content,
744            total_size,
745            event_count: accumulator.events.len(),
746            has_message_start: accumulator.has_message_start,
747            sse_events: sse_events_json,
748        }
749        .to_event(original_event)
750    }
751
752    fn merged_json_content(accumulator: &SSEAccumulator) -> String {
753        let has_openai_json =
754            !accumulator.openai_reasoning.is_empty() || !accumulator.openai_tool_calls.is_empty();
755        if !has_openai_json {
756            return Self::formatted_accumulated_json(&accumulator.accumulated_json);
757        }
758
759        let mut merged = serde_json::Map::new();
760        if !accumulator.accumulated_json.is_empty() {
761            match serde_json::from_str::<Value>(&accumulator.accumulated_json) {
762                Ok(parsed_json) => {
763                    merged.insert("partial_json".to_string(), parsed_json);
764                }
765                Err(_) => {
766                    merged.insert(
767                        "partial_json".to_string(),
768                        Value::String(accumulator.accumulated_json.clone()),
769                    );
770                }
771            }
772        }
773        if !accumulator.openai_reasoning.is_empty() {
774            merged.insert(
775                "reasoning_content".to_string(),
776                Value::String(accumulator.openai_reasoning.clone()),
777            );
778        }
779        if !accumulator.openai_tool_calls.is_empty() {
780            let tool_calls = accumulator
781                .openai_tool_calls
782                .iter()
783                .map(|(index, tool_call)| {
784                    json!({
785                        "index": index,
786                        "id": tool_call.id,
787                        "type": tool_call.type_name,
788                        "function": {
789                            "name": tool_call.function_name,
790                            "arguments": tool_call.function_arguments,
791                        }
792                    })
793                })
794                .collect::<Vec<_>>();
795            merged.insert("tool_calls".to_string(), Value::Array(tool_calls));
796        }
797        if merged.is_empty() {
798            String::new()
799        } else {
800            serde_json::to_string_pretty(&Value::Object(merged)).unwrap_or_default()
801        }
802    }
803
804    fn formatted_accumulated_json(json_content: &str) -> String {
805        if json_content.is_empty() {
806            String::new()
807        } else if let Ok(parsed_json) = serde_json::from_str::<Value>(json_content) {
808            serde_json::to_string_pretty(&parsed_json).unwrap_or_else(|_| json_content.to_string())
809        } else {
810            json_content.to_string()
811        }
812    }
813
814    fn event_host(event: &Event) -> Option<String> {
815        event
816            .data
817            .get("host")
818            .and_then(|v| v.as_str())
819            .or_else(|| {
820                event
821                    .data
822                    .get("headers")
823                    .and_then(|headers| headers.as_object())
824                    .and_then(|headers| {
825                        headers.iter().find_map(|(key, value)| {
826                            (key.eq_ignore_ascii_case("host") || key == ":authority")
827                                .then(|| value.as_str())
828                                .flatten()
829                        })
830                    })
831            })
832            .map(str::to_string)
833    }
834
835    fn evict_over_capacity(buffers: &mut HashMap<String, SSEAccumulator>, max: usize) {
836        while buffers.len() > max {
837            let oldest_key = buffers
838                .iter()
839                .min_by_key(|(_, acc)| acc.last_update)
840                .map(|(k, _)| k.clone());
841            if let Some(key) = oldest_key {
842                buffers.remove(&key);
843            } else {
844                break;
845            }
846        }
847    }
848}
849
850#[async_trait]
851impl Analyzer for SSEProcessor {
852    async fn process(&mut self, stream: EventStream) -> Result<EventStream, AnalyzerError> {
853        let sse_buffers = Arc::clone(&self.sse_buffers);
854        let timeout_ms = self.timeout_ms;
855        let max_buffers = self.max_buffers;
856
857        let processed_stream = stream.filter_map(move |event| {
858            let buffers = Arc::clone(&sse_buffers);
859
860            async move {
861                let Some((data_str, allow_json_fragment)) = Self::sse_payload(&event) else {
862                    return Some(event);
863                };
864
865                let sse_events = if Self::is_sse_data(data_str) {
866                    Self::parse_sse_events(data_str)
867                } else if allow_json_fragment
868                    && let Some(event) = Self::parse_usage_metadata_fragment(data_str)
869                {
870                    vec![event]
871                } else {
872                    return Some(event);
873                };
874                if sse_events.is_empty() {
875                    return Some(event);
876                }
877
878                let has_content_potential = sse_events.iter().any(|sse_event| {
879                    if let Some(event_type) = &sse_event.event {
880                        !matches!(event_type.as_str(), "message_delta" | "ping")
881                    } else {
882                        true
883                    }
884                });
885
886                let should_skip_chunk = !has_content_potential
887                    && sse_events.iter().all(|e| {
888                        e.event
889                            .as_deref()
890                            .is_some_and(|t| matches!(t, "ping" | "message_delta"))
891                    });
892
893                if should_skip_chunk {
894                    let connection_id = Self::generate_connection_id(&event, &sse_events);
895                    let buffers_lock = buffers.lock().unwrap();
896                    let has_existing = buffers_lock.contains_key(&connection_id);
897                    drop(buffers_lock);
898                    if !has_existing {
899                        return None;
900                    }
901                }
902
903                let connection_id = Self::generate_connection_id(&event, &sse_events);
904
905                let mut buffers_lock = buffers.lock().unwrap();
906
907                buffers_lock
908                    .retain(|_, acc| event.timestamp.saturating_sub(acc.last_update) <= timeout_ms);
909                Self::evict_over_capacity(&mut buffers_lock, max_buffers);
910
911                let mut final_connection_id = connection_id.clone();
912
913                if let Some(message_id) = Self::extract_message_id(&sse_events) {
914                    let pid = event.data.get("pid").and_then(|v| v.as_u64()).unwrap_or(0);
915                    let tid = event.data.get("tid").and_then(|v| v.as_u64()).unwrap_or(0);
916                    final_connection_id = format!("{}:{}:{}", pid, tid, message_id);
917                } else {
918                    let pid = event.data.get("pid").and_then(|v| v.as_u64()).unwrap_or(0);
919                    let tid = event.data.get("tid").and_then(|v| v.as_u64()).unwrap_or(0);
920                    let conn_prefix = format!("{}:{}:", pid, tid);
921
922                    for (existing_id, accumulator) in buffers_lock.iter() {
923                        if existing_id.starts_with(&conn_prefix) && !accumulator.is_complete {
924                            let has_message_stop = accumulator
925                                .events
926                                .iter()
927                                .any(|e| e.event.as_deref() == Some("message_stop"));
928                            if !has_message_stop {
929                                final_connection_id = existing_id.clone();
930                                break;
931                            }
932                        }
933                    }
934                }
935
936                let accumulator = buffers_lock
937                    .entry(final_connection_id.clone())
938                    .or_insert_with(|| SSEAccumulator {
939                        message_id: None,
940                        accumulated_text: String::new(),
941                        accumulated_json: String::new(),
942                        openai_reasoning: String::new(),
943                        openai_tool_calls: BTreeMap::new(),
944                        events: Vec::new(),
945                        is_complete: false,
946                        last_update: event.timestamp,
947                        has_message_start: false,
948                        start_time: event.timestamp,
949                        end_time: event.timestamp,
950                    });
951
952                accumulator.last_update = event.timestamp;
953                accumulator.end_time = event.timestamp;
954
955                Self::accumulate_content(accumulator, &sse_events);
956
957                let terminal_finish_completes_http_body = !allow_json_fragment
958                    && sse_events.iter().any(Self::sse_event_has_terminal_finish);
959
960                if Self::is_sse_complete(accumulator) || terminal_finish_completes_http_body {
961                    let result_event = if Self::has_meaningful_content(accumulator) {
962                        Some(Self::create_merged_event(
963                            final_connection_id.clone(),
964                            accumulator,
965                            &event,
966                        ))
967                    } else {
968                        None
969                    };
970
971                    buffers_lock.remove(&final_connection_id);
972                    drop(buffers_lock);
973
974                    result_event
975                } else {
976                    None
977                }
978            }
979        });
980
981        Ok(Box::pin(processed_stream))
982    }
983}
984
985fn extract_json_object_after_key<'a>(text: &'a str, key: &str) -> Option<&'a str> {
986    let key_index = text.find(key)?;
987    let object_start = text[key_index..].find('{')? + key_index;
988    let mut depth = 0usize;
989    let mut in_string = false;
990    let mut escape = false;
991
992    for (offset, ch) in text[object_start..].char_indices() {
993        if in_string {
994            if escape {
995                escape = false;
996            } else if ch == '\\' {
997                escape = true;
998            } else if ch == '"' {
999                in_string = false;
1000            }
1001            continue;
1002        }
1003
1004        match ch {
1005            '"' => in_string = true,
1006            '{' => depth += 1,
1007            '}' => {
1008                depth = depth.saturating_sub(1);
1009                if depth == 0 {
1010                    let end = object_start + offset + ch.len_utf8();
1011                    return Some(&text[object_start..end]);
1012                }
1013            }
1014            _ => {}
1015        }
1016    }
1017
1018    None
1019}
1020
1021fn extract_json_string_field(text: &str, key: &str) -> Option<String> {
1022    let key_pattern = format!("\"{}\"", key);
1023    let key_index = text.find(&key_pattern)?;
1024    let after_key = &text[key_index + key_pattern.len()..];
1025    let colon = after_key.find(':')?;
1026    let mut chars = after_key[colon + 1..].char_indices().peekable();
1027    while let Some((_, ch)) = chars.peek().copied() {
1028        if ch.is_whitespace() {
1029            chars.next();
1030        } else {
1031            break;
1032        }
1033    }
1034    let (start_offset, quote) = chars.next()?;
1035    if quote != '"' {
1036        return None;
1037    }
1038    let value_start = key_index + key_pattern.len() + colon + 1 + start_offset + quote.len_utf8();
1039    let rest = &text[value_start..];
1040    let mut escape = false;
1041    for (offset, ch) in rest.char_indices() {
1042        if escape {
1043            escape = false;
1044        } else if ch == '\\' {
1045            escape = true;
1046        } else if ch == '"' {
1047            let raw = &rest[..offset];
1048            return serde_json::from_str::<String>(&format!("\"{}\"", raw)).ok();
1049        }
1050    }
1051    None
1052}