gproxy_transform/transform/stream_adapter/
buffered.rs1use super::{SseDecoder, SseTransformer};
2use crate::protocol::ContentGenerationKind;
3
4pub fn convert_buffered(mut transformer: SseTransformer, body: &[u8]) -> Vec<u8> {
6 let mut out = transformer.push(body);
7 out.extend(transformer.finish());
8 out
9}
10
11pub fn aggregate_buffered(kind: ContentGenerationKind, sse_body: &[u8]) -> Vec<u8> {
13 use crate::transform::generate_content::stream_to_response as s2r;
14 use ContentGenerationKind as K;
15
16 let mut decoder = SseDecoder::new();
17 let mut frames = decoder.push(sse_body);
18 if let Some(tail) = decoder.finish() {
19 frames.push(tail);
20 }
21 let datas: Vec<String> = frames
22 .into_iter()
23 .map(|frame| frame.data)
24 .filter(|data| data.trim() != "[DONE]")
25 .collect();
26
27 macro_rules! collapse {
28 ($ty:ty, $aggregate:path) => {{
29 let events = datas
30 .iter()
31 .filter_map(|data| serde_json::from_str::<$ty>(data).ok());
32 serde_json::to_vec(&$aggregate(events))
33 }};
34 }
35
36 let out = match kind {
37 K::OpenAiResponses | K::OpenAiResponsesWebSocket => collapse!(
38 crate::protocol::openai::ResponseStreamEvent,
39 s2r::openai_responses::response
40 ),
41 K::OpenAiChatCompletions => collapse!(
42 crate::protocol::openai::ChatCompletionChunk,
43 s2r::openai_chat::response
44 ),
45 K::ClaudeMessages => collapse!(
46 crate::protocol::claude::StreamEvent,
47 s2r::claude_messages::response
48 ),
49 K::GeminiGenerateContent => collapse!(
50 crate::protocol::gemini::StreamGenerateContentChunk,
51 s2r::gemini_generate_content::response
52 ),
53 };
54 out.unwrap_or_else(|_| sse_body.to_vec())
55}