gproxy_channel_api/
transport.rs1use bytes::Bytes;
4
5#[derive(Debug, Clone, Copy, Default)]
21pub struct TransportOptions {
22 pub total_timeout: Option<std::time::Duration>,
23 pub max_redirects: Option<usize>,
24 pub http_version: Option<http::Version>,
25 pub omit_body: bool,
26}
27
28pub trait ByteStreamDecoder: Send {
30 fn push(&mut self, chunk: &[u8]) -> Vec<u8>;
32
33 fn finish(&mut self) -> Vec<u8>;
35}
36
37#[derive(Debug, thiserror::Error)]
39pub enum ClientError {
40 #[error("upstream transport error: {0}")]
41 Transport(String),
42 #[error("upstream client config error: {0}")]
45 Config(String),
46}
47
48#[cfg(not(target_arch = "wasm32"))]
51pub type RespStream =
52 std::pin::Pin<Box<dyn futures_core::Stream<Item = Result<Bytes, ClientError>> + Send>>;
53#[cfg(target_arch = "wasm32")]
54pub type RespStream =
55 std::pin::Pin<Box<dyn futures_core::Stream<Item = Result<Bytes, ClientError>>>>;
56
57#[cfg(not(target_arch = "wasm32"))]
59#[derive(Debug)]
60pub enum ConduitFrame {
61 Text(String),
62 Binary(Bytes),
63 Close,
64}
65
66#[cfg(not(target_arch = "wasm32"))]
68#[async_trait::async_trait]
69pub trait ConduitSocket: Send {
70 async fn send_text(&mut self, text: String) -> Result<(), ClientError>;
71 async fn send_binary(&mut self, _bytes: Bytes) -> Result<(), ClientError> {
72 Err(ClientError::Config(
73 "binary upstream websocket frames not supported by this client".into(),
74 ))
75 }
76 async fn recv_text(&mut self) -> Option<Result<String, ClientError>>;
77 async fn recv_frame(&mut self) -> Option<Result<ConduitFrame, ClientError>> {
78 self.recv_text()
79 .await
80 .map(|result| result.map(ConduitFrame::Text))
81 }
82 async fn close(&mut self) -> Result<(), ClientError> {
83 Ok(())
84 }
85}
86
87#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
91#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
92pub trait UpstreamClient: Send + Sync {
93 async fn send(&self, req: http::Request<Bytes>) -> Result<http::Response<Bytes>, ClientError>;
94
95 async fn send_websocket(
96 &self,
97 _req: http::Request<Bytes>,
98 ) -> Result<http::Response<Bytes>, ClientError> {
99 Err(ClientError::Config(
100 "upstream websocket not supported by this client".into(),
101 ))
102 }
103
104 async fn send_streaming(
105 &self,
106 req: http::Request<Bytes>,
107 ) -> Result<(http::StatusCode, http::HeaderMap, RespStream), ClientError> {
108 let resp = self.send(req).await?;
109 let (parts, body) = resp.into_parts();
110 let once = futures_util::stream::once(async move { Ok::<Bytes, ClientError>(body) });
111 Ok((parts.status, parts.headers, Box::pin(once)))
112 }
113
114 #[cfg(not(target_arch = "wasm32"))]
115 async fn open_websocket(
116 &self,
117 _req: http::Request<Bytes>,
118 ) -> Result<Box<dyn ConduitSocket>, ClientError> {
119 Err(ClientError::Config(
120 "upstream websocket not supported by this client".into(),
121 ))
122 }
123}