Skip to main content

gproxy_transform/envelope/response/
mod.rs

1mod framing;
2
3use bytes::Bytes;
4use gproxy_protocol::{OperationKey, StreamFraming};
5
6use self::framing::{FrameDecoder, FrameEncoder};
7use super::SseFrame;
8use crate::TransformError;
9use crate::registry::{self, TransformPair};
10
11pub struct ResponseStream {
12    decoder: FrameDecoder,
13    converter: Box<dyn Converter>,
14    encoder: FrameEncoder,
15}
16
17pub(crate) trait Converter: Send {
18    fn frame(&mut self, frame: SseFrame) -> Result<Vec<Bytes>, TransformError>;
19    fn finish(&mut self) -> Result<Vec<Bytes>, TransformError>;
20}
21
22impl ResponseStream {
23    pub fn new(source: OperationKey, target: OperationKey) -> Result<Self, TransformError> {
24        Self::new_framed(source, target, StreamFraming::Sse, StreamFraming::Sse)
25    }
26
27    pub fn new_framed(
28        source: OperationKey,
29        target: OperationKey,
30        source_framing: StreamFraming,
31        target_framing: StreamFraming,
32    ) -> Result<Self, TransformError> {
33        let converter: Box<dyn Converter> = if source == target {
34            Box::new(Passthrough)
35        } else {
36            let pair =
37                registry::resolve(source, target).ok_or(TransformError::UnsupportedPair {
38                    source_key: source,
39                    target_key: target,
40                })?;
41            converter(pair, source, target)?
42        };
43        Ok(Self {
44            decoder: FrameDecoder::new(target_framing)?,
45            converter,
46            encoder: FrameEncoder::new(source_framing)?,
47        })
48    }
49
50    pub fn push(&mut self, chunk: Bytes) -> Result<Vec<Bytes>, TransformError> {
51        let frames = self.decoder.push(&chunk)?;
52        self.convert(frames)
53    }
54
55    pub fn finish(&mut self) -> Result<Vec<Bytes>, TransformError> {
56        let frames = self.decoder.finish()?;
57        let mut output = self.convert(frames)?;
58        let terminal = self.converter.finish()?;
59        output.extend(self.encoder.push(terminal)?);
60        output.extend(self.encoder.finish()?);
61        Ok(output)
62    }
63
64    fn convert(&mut self, frames: Vec<SseFrame>) -> Result<Vec<Bytes>, TransformError> {
65        let mut output = Vec::new();
66        for frame in frames {
67            output.extend(self.encoder.push(self.converter.frame(frame)?)?);
68        }
69        Ok(output)
70    }
71}
72
73struct Passthrough;
74
75impl Converter for Passthrough {
76    fn frame(&mut self, frame: SseFrame) -> Result<Vec<Bytes>, TransformError> {
77        Ok(vec![SseFrame::encode(frame._event.as_deref(), &frame.data)])
78    }
79
80    fn finish(&mut self) -> Result<Vec<Bytes>, TransformError> {
81        Ok(Vec::new())
82    }
83}
84
85fn converter(
86    pair: TransformPair,
87    source: OperationKey,
88    target: OperationKey,
89) -> Result<Box<dyn Converter>, TransformError> {
90    Ok(match pair {
91        TransformPair::ChatToClaude => {
92            crate::generate_content::openai_chat_to_claude_messages::stream::converter()
93        }
94        TransformPair::ResponsesToClaude => {
95            crate::generate_content::openai_responses_to_claude_messages::stream::converter()
96        }
97        TransformPair::ClaudeToChat => {
98            crate::generate_content::claude_messages_to_openai_chat::stream::converter()
99        }
100        TransformPair::ClaudeToResponses => {
101            crate::generate_content::claude_messages_to_openai_responses::stream::converter()
102        }
103        TransformPair::ClaudeToGemini => {
104            crate::generate_content::gemini_generate_content_to_claude_messages::stream::converter()
105        }
106        TransformPair::GeminiToClaude => {
107            crate::generate_content::claude_messages_to_gemini_generate_content::stream::converter()
108        }
109        TransformPair::GeminiToChat => {
110            crate::generate_content::gemini_generate_content_to_openai_chat::stream::converter()
111        }
112        TransformPair::ChatToGemini => {
113            crate::generate_content::openai_chat_to_gemini_generate_content::stream::converter()
114        }
115        TransformPair::GeminiToResponses => {
116            crate::generate_content::gemini_generate_content_to_openai_responses::stream::converter(
117            )
118        }
119        TransformPair::ResponsesToGemini => {
120            crate::generate_content::openai_responses_to_gemini_generate_content::stream::converter(
121            )
122        }
123        TransformPair::OpenAiChatToResponses => {
124            crate::generate_content::openai_chat_to_openai_responses::stream::converter()
125        }
126        TransformPair::OpenAiResponsesToChat => {
127            crate::generate_content::openai_responses_to_openai_chat::stream::converter()
128        }
129        TransformPair::OpenAiCreateImageToGemini => crate::images::stream::from_gemini(false),
130        TransformPair::OpenAiEditImageToGemini => crate::images::stream::from_gemini(true),
131        TransformPair::OpenAiCreateImageToResponses => crate::images::stream::from_responses(false),
132        TransformPair::OpenAiEditImageToResponses => crate::images::stream::from_responses(true),
133        _ => {
134            return Err(TransformError::UnsupportedPair {
135                source_key: source,
136                target_key: target,
137            });
138        }
139    })
140}