Skip to main content

fusen_rs/protocol/codec/
mod.rs

1use crate::{
2    error::FusenError,
3    protocol::{
4        codec::body::{RequestBodyCodec, ResponseBodyCodec, json::JsonCodec, triple::TripleCodec},
5        fusen::{
6            request::{FusenRequest, Path},
7            response::{FusenResponse, HttpStatus},
8        },
9    },
10};
11use bytes::{Bytes, BytesMut};
12use fusen_internal_common::protocol::Protocol;
13use http::{
14    Request, Response, Version,
15    header::{CONNECTION, CONTENT_TYPE},
16};
17use http_body_util::{BodyExt, Full, combinators::BoxBody};
18use std::{collections::HashMap, convert::Infallible};
19
20pub mod body;
21
22#[derive(Default)]
23pub struct FusenHttpCodec {
24    json_codec: JsonCodec,
25    triple_codec: TripleCodec,
26}
27
28impl RequestCodec<Bytes, hyper::Error> for FusenHttpCodec {
29    fn encode(
30        &self,
31        fusen_request: &mut FusenRequest,
32    ) -> Result<
33        Request<http_body_util::combinators::BoxBody<Bytes, std::convert::Infallible>>,
34        crate::error::FusenError,
35    > {
36        let mut builder = Request::builder().header(CONNECTION, "keep-alive");
37        for (key, value) in fusen_request.headers.drain() {
38            builder = builder.header(key, value);
39        }
40        let Some(addr) = &fusen_request.addr else {
41            return Err(FusenError::Impossible);
42        };
43        let mut uri = format!("{}{}", addr, fusen_request.path.path);
44        if !fusen_request.querys.is_empty() {
45            if fusen_request.path.path.contains('{') {
46                let mut path = fusen_request.path.path.to_string();
47                for (key, value) in &fusen_request.querys {
48                    path = path.replace(&format!("{{{key}}}"), value);
49                }
50                uri = format!("{addr}{path}");
51            } else {
52                uri.push('?');
53                for (key, value) in &fusen_request.querys {
54                    uri.push_str(&format!("{}={}&", key, urlencoding::encode(value.as_str())));
55                }
56                uri.pop();
57            }
58        }
59        let mut body = Bytes::new();
60        let mut version = Version::HTTP_2;
61        if let Protocol::Host(_) | Protocol::SpringCloud(_) = fusen_request.protocol {
62            version = Version::HTTP_11;
63        }
64        if let Some(bodys) = fusen_request.bodys.take() {
65            match &fusen_request.protocol {
66                Protocol::Dubbo => {
67                    builder = builder.header(CONTENT_TYPE, "application/grpc");
68                    body = RequestBodyCodec::encode(&self.triple_codec, bodys)?;
69                }
70                _ => {
71                    builder = builder.header(CONTENT_TYPE, "application/json");
72                    body = RequestBodyCodec::encode(&self.json_codec, bodys)?;
73                }
74            };
75        }
76        builder
77            .version(version)
78            .method(fusen_request.path.method.clone())
79            .uri(uri)
80            .body(Full::new(body).boxed())
81            .map_err(|error| crate::error::FusenError::Error(Box::new(error)))
82    }
83
84    async fn decode(
85        &self,
86        mut request: Request<http_body_util::combinators::BoxBody<Bytes, hyper::Error>>,
87    ) -> Result<FusenRequest, crate::error::FusenError> {
88        let mut querys = HashMap::new();
89        let mut headers = HashMap::new();
90        for (key, value) in request.headers_mut().drain() {
91            let Some(key) = key else {
92                continue;
93            };
94            headers.insert(
95                key.to_string().to_ascii_lowercase(),
96                String::from_utf8_lossy(value.as_bytes()).to_string(),
97            );
98        }
99        if let Some(request_querys) = request.uri().query() {
100            let request_querys: Vec<&str> = request_querys.split('&').collect();
101            for query in request_querys {
102                if let Some((key, value)) = query.split_once('=') {
103                    querys.insert(
104                        key.to_string(),
105                        urlencoding::decode(value)
106                            .map(|e| e.to_string())
107                            .unwrap_or(value.to_string()),
108                    );
109                }
110            }
111        }
112        let mut protocol = Protocol::Fusen;
113        let mut bodys = None;
114        if let Some(content_type) = headers.get(CONTENT_TYPE.as_str()) {
115            let bytes = read_body(request.body_mut()).await.freeze();
116            if content_type.starts_with("application/grpc") {
117                protocol = Protocol::Dubbo;
118                let _ = bodys.insert(RequestBodyCodec::decode(&self.triple_codec, bytes)?);
119            } else if content_type.starts_with("application/json") {
120                let _ = bodys.insert(RequestBodyCodec::decode(&self.json_codec, bytes)?);
121            }
122        };
123        Ok(FusenRequest {
124            path: Path {
125                method: request.method().clone(),
126                path: request.uri().path().to_owned(),
127            },
128            addr: None,
129            querys,
130            headers,
131            extensions: None,
132            bodys,
133            protocol,
134        })
135    }
136}
137
138impl ResponseCodec<Bytes, hyper::Error> for FusenHttpCodec {
139    fn encode(
140        &self,
141        fusen_response: &mut FusenResponse,
142    ) -> Result<
143        http::Response<http_body_util::combinators::BoxBody<Bytes, std::convert::Infallible>>,
144        crate::error::FusenError,
145    > {
146        let mut builder = Response::builder();
147        for (key, value) in fusen_response.headers.drain() {
148            builder = builder.header(key, value);
149        }
150        let mut body = Bytes::new();
151        if let Some(bodys) = fusen_response.body.take() {
152            match &fusen_response.protocol {
153                Protocol::Dubbo => {
154                    builder = builder.header(CONTENT_TYPE, "application/grpc");
155                    body = ResponseBodyCodec::encode(&self.triple_codec, bodys)?;
156                }
157                _ => {
158                    builder = builder.header(CONTENT_TYPE, "application/json");
159                    body = ResponseBodyCodec::encode(&self.json_codec, bodys)?;
160                }
161            };
162        }
163        builder
164            .status(fusen_response.http_status.status)
165            .body(Full::new(body).boxed())
166            .map_err(|error| FusenError::Error(Box::new(error)))
167    }
168
169    async fn decode(
170        &self,
171        mut response: http::Response<http_body_util::combinators::BoxBody<Bytes, hyper::Error>>,
172    ) -> Result<FusenResponse, crate::error::FusenError> {
173        let mut headers: HashMap<String, String> = HashMap::new();
174        for (key, value) in response.headers_mut().drain() {
175            let Some(key) = key else {
176                continue;
177            };
178            headers.insert(
179                key.to_string().to_ascii_lowercase(),
180                String::from_utf8_lossy(value.as_bytes()).to_string(),
181            );
182        }
183        let mut protocol = Protocol::Fusen;
184        let mut body = None;
185        if let Some(content_type) = headers.get(CONTENT_TYPE.as_str()) {
186            let bytes = read_body(response.body_mut()).await.freeze();
187            if content_type.starts_with("application/grpc") {
188                protocol = Protocol::Dubbo;
189                let _ = body.insert(ResponseBodyCodec::decode(&self.triple_codec, bytes)?);
190            } else if content_type.starts_with("application/json") {
191                let _ = body.insert(ResponseBodyCodec::decode(&self.json_codec, bytes)?);
192            }
193        };
194        Ok(FusenResponse {
195            protocol,
196            http_status: HttpStatus {
197                status: response.status().as_u16(),
198                message: None,
199            },
200            headers,
201            extensions: None,
202            body,
203        })
204    }
205}
206
207#[allow(async_fn_in_trait)]
208pub trait RequestCodec<T, E> {
209    fn encode(
210        &self,
211        fusen_request: &mut FusenRequest,
212    ) -> Result<Request<BoxBody<T, Infallible>>, FusenError>;
213
214    async fn decode(&self, request: Request<BoxBody<T, E>>) -> Result<FusenRequest, FusenError>;
215}
216
217#[allow(async_fn_in_trait)]
218pub trait ResponseCodec<T, E> {
219    fn encode(
220        &self,
221        fusen_response: &mut FusenResponse,
222    ) -> Result<Response<BoxBody<T, Infallible>>, FusenError>;
223
224    async fn decode(&self, response: Response<BoxBody<T, E>>) -> Result<FusenResponse, FusenError>;
225}
226
227async fn read_body(body: &mut BoxBody<Bytes, hyper::Error>) -> BytesMut {
228    let mut mut_bytes = BytesMut::new();
229    while let Some(Ok(frame)) = body.frame().await {
230        if let Ok(bytes) = frame.into_data() {
231            mut_bytes.extend_from_slice(&bytes);
232        }
233    }
234    mut_bytes
235}