1use 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}