gproxy_transform/envelope/response/
mod.rs1mod 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}