This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / bobbin / crates / runtime / src / network.rs
10 kB 307 lines
1use std::future::Future; 2use std::net::SocketAddr; 3use std::pin::Pin; 4use std::sync::Arc; 5 6use bytes::Bytes; 7use futures::stream::{Stream, StreamExt}; 8use http::{HeaderMap, StatusCode}; 9use thiserror::Error; 10use tokio_tungstenite::tungstenite::{ 11 Bytes as WsBytes, Message as TungsteniteMessage, protocol::CloseFrame as TungsteniteClose, 12 protocol::frame::coding::CloseCode as TungsteniteCloseCode, 13}; 14use url::Url; 15 16#[derive(Debug, Error)] 17pub enum NetworkError { 18 #[error("connect: {0}")] 19 Connect(String), 20 #[error("timeout: {0}")] 21 Timeout(String), 22 #[error("redirect: {0}")] 23 Redirect(String), 24 #[error("transport: {0}")] 25 Transport(String), 26 #[error("body: {0}")] 27 Body(String), 28 #[error("protocol: {0}")] 29 Protocol(String), 30} 31 32pub struct HttpRequest { 33 pub url: Url, 34 pub headers: HeaderMap, 35} 36 37pub type BodyStream = Pin<Box<dyn Stream<Item = Result<Bytes, NetworkError>> + Send + 'static>>; 38 39pub struct HttpResponseHead { 40 pub status: StatusCode, 41 pub headers: HeaderMap, 42 pub content_length: Option<u64>, 43 pub body: BodyStream, 44} 45 46pub type HttpResult = Result<HttpResponseHead, NetworkError>; 47pub type HttpResponseFuture = Pin<Box<dyn Future<Output = HttpResult> + Send + 'static>>; 48 49pub trait HttpTransport: Send + Sync + 'static { 50 fn execute(&self, request: HttpRequest) -> HttpResponseFuture; 51} 52 53#[derive(Clone, Debug)] 54pub struct ReqwestHttp { 55 client: reqwest::Client, 56} 57 58impl ReqwestHttp { 59 pub fn new(client: reqwest::Client) -> Self { 60 Self { client } 61 } 62 63 pub fn shared(client: reqwest::Client) -> Arc<dyn HttpTransport> { 64 Arc::new(Self::new(client)) 65 } 66} 67 68impl HttpTransport for ReqwestHttp { 69 fn execute(&self, request: HttpRequest) -> HttpResponseFuture { 70 let client = self.client.clone(); 71 Box::pin(async move { 72 let resp = client 73 .get(request.url) 74 .headers(request.headers) 75 .send() 76 .await 77 .map_err(map_reqwest)?; 78 let status = resp.status(); 79 let headers = resp.headers().clone(); 80 let content_length = resp.content_length(); 81 let body: BodyStream = Box::pin( 82 resp.bytes_stream() 83 .map(|chunk| chunk.map_err(|e| NetworkError::Body(e.to_string()))), 84 ); 85 Ok(HttpResponseHead { 86 status, 87 headers, 88 content_length, 89 body, 90 }) 91 }) 92 } 93} 94 95/// Lets jacquard resolve identities over the workspace reqwest (0.13); jacquard's own 96/// `HttpClient` impl is against reqwest 0.12, which is built here without TLS. 97impl jacquard_common::http_client::HttpClient for ReqwestHttp { 98 type Error = reqwest::Error; 99 100 async fn send_http( 101 &self, 102 request: http::Request<Vec<u8>>, 103 ) -> Result<http::Response<Vec<u8>>, reqwest::Error> { 104 let (parts, body) = request.into_parts(); 105 let mut req = self 106 .client 107 .request(parts.method, parts.uri.to_string()) 108 .body(body); 109 for (name, value) in parts.headers.iter() { 110 req = req.header(name, value); 111 } 112 113 let resp = req.send().await?; 114 let mut builder = http::Response::builder().status(resp.status()); 115 for (name, value) in resp.headers().iter() { 116 builder = builder.header(name, value); 117 } 118 let body = resp.bytes().await?.to_vec(); 119 Ok(builder.body(body).expect("response parts came from reqwest")) 120 } 121} 122 123fn map_reqwest(err: reqwest::Error) -> NetworkError { 124 let msg = err.to_string(); 125 if err.is_timeout() { 126 NetworkError::Timeout(msg) 127 } else if err.is_connect() { 128 NetworkError::Connect(msg) 129 } else if err.is_redirect() { 130 NetworkError::Redirect(msg) 131 } else { 132 NetworkError::Transport(msg) 133 } 134} 135 136#[derive(Clone, Debug)] 137pub enum WsMessage { 138 Text(String), 139 Binary(Bytes), 140 Ping(Bytes), 141 Pong(Bytes), 142 Close { code: u16, reason: String }, 143} 144 145pub type WsSendFuture<'a> = Pin<Box<dyn Future<Output = Result<(), NetworkError>> + Send + 'a>>; 146pub type WsMessageFuture<'a> = 147 Pin<Box<dyn Future<Output = Option<Result<WsMessage, NetworkError>>> + Send + 'a>>; 148 149pub trait WsSink: Send + 'static { 150 fn send<'a>(&'a mut self, message: WsMessage) -> WsSendFuture<'a>; 151} 152 153pub trait WsStream: Send + 'static { 154 fn next<'a>(&'a mut self) -> WsMessageFuture<'a>; 155} 156 157pub struct WsConn { 158 pub sink: Box<dyn WsSink>, 159 pub stream: Box<dyn WsStream>, 160} 161 162pub type WsConnectFuture = 163 Pin<Box<dyn Future<Output = Result<WsConn, NetworkError>> + Send + 'static>>; 164 165pub trait WsTransport: Send + Sync + 'static { 166 fn connect(&self, url: Url) -> WsConnectFuture; 167} 168 169pub type AddrGuard = Arc<dyn Fn(&[SocketAddr]) -> Result<(), NetworkError> + Send + Sync>; 170 171#[derive(Clone, Copy, Debug, Default)] 172pub struct TungsteniteWs; 173 174impl TungsteniteWs { 175 pub fn shared() -> Arc<dyn WsTransport> { 176 Arc::new(Self) 177 } 178} 179 180impl WsTransport for TungsteniteWs { 181 fn connect(&self, url: Url) -> WsConnectFuture { 182 Box::pin(async move { 183 let url_str = url.as_str().to_owned(); 184 let (ws, _resp) = tokio_tungstenite::connect_async(&url_str) 185 .await 186 .map_err(|e| NetworkError::Connect(e.to_string()))?; 187 let (sink_inner, stream_inner) = futures::StreamExt::split(ws); 188 let sink: Box<dyn WsSink> = Box::new(TungsteniteSink { inner: sink_inner }); 189 let stream: Box<dyn WsStream> = Box::new(TungsteniteStream { 190 inner: stream_inner, 191 }); 192 Ok(WsConn { sink, stream }) 193 }) 194 } 195} 196 197pub struct GuardedWs { 198 guard: AddrGuard, 199} 200 201impl GuardedWs { 202 pub fn shared(guard: AddrGuard) -> Arc<dyn WsTransport> { 203 Arc::new(Self { guard }) 204 } 205} 206 207impl WsTransport for GuardedWs { 208 fn connect(&self, url: Url) -> WsConnectFuture { 209 let guard = self.guard.clone(); 210 Box::pin(async move { 211 let host = url 212 .host_str() 213 .ok_or_else(|| NetworkError::Connect("ws url missing host".to_owned()))? 214 .to_owned(); 215 let port = url 216 .port_or_known_default() 217 .ok_or_else(|| NetworkError::Connect("ws url missing port".to_owned()))?; 218 let addrs: Vec<SocketAddr> = tokio::net::lookup_host((host.as_str(), port)) 219 .await 220 .map_err(|e| NetworkError::Connect(e.to_string()))? 221 .collect(); 222 guard(&addrs)?; 223 let addr = addrs 224 .into_iter() 225 .next() 226 .ok_or_else(|| NetworkError::Connect(format!("no addresses for {host}")))?; 227 let tcp = tokio::net::TcpStream::connect(addr) 228 .await 229 .map_err(|e| NetworkError::Connect(e.to_string()))?; 230 let (ws, _resp) = tokio_tungstenite::client_async_tls(url.as_str(), tcp) 231 .await 232 .map_err(|e| NetworkError::Connect(e.to_string()))?; 233 let (sink_inner, stream_inner) = futures::StreamExt::split(ws); 234 let sink: Box<dyn WsSink> = Box::new(TungsteniteSink { inner: sink_inner }); 235 let stream: Box<dyn WsStream> = Box::new(TungsteniteStream { 236 inner: stream_inner, 237 }); 238 Ok(WsConn { sink, stream }) 239 }) 240 } 241} 242 243type TungsteniteWsStream = 244 tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>; 245 246struct TungsteniteSink { 247 inner: futures::stream::SplitSink<TungsteniteWsStream, TungsteniteMessage>, 248} 249 250impl WsSink for TungsteniteSink { 251 fn send<'a>(&'a mut self, message: WsMessage) -> WsSendFuture<'a> { 252 Box::pin(async move { 253 use futures::SinkExt; 254 self.inner 255 .send(message_to_tungstenite(message)) 256 .await 257 .map_err(|e| NetworkError::Transport(e.to_string())) 258 }) 259 } 260} 261 262struct TungsteniteStream { 263 inner: futures::stream::SplitStream<TungsteniteWsStream>, 264} 265 266impl WsStream for TungsteniteStream { 267 fn next<'a>(&'a mut self) -> WsMessageFuture<'a> { 268 Box::pin(async move { 269 let item = StreamExt::next(&mut self.inner).await?; 270 Some( 271 item.map_err(|e| NetworkError::Transport(e.to_string())) 272 .and_then(message_from_tungstenite), 273 ) 274 }) 275 } 276} 277 278fn message_to_tungstenite(message: WsMessage) -> TungsteniteMessage { 279 match message { 280 WsMessage::Text(text) => TungsteniteMessage::Text(text.into()), 281 WsMessage::Binary(bytes) => TungsteniteMessage::Binary(WsBytes::copy_from_slice(&bytes)), 282 WsMessage::Ping(bytes) => TungsteniteMessage::Ping(WsBytes::copy_from_slice(&bytes)), 283 WsMessage::Pong(bytes) => TungsteniteMessage::Pong(WsBytes::copy_from_slice(&bytes)), 284 WsMessage::Close { code, reason } => TungsteniteMessage::Close(Some(TungsteniteClose { 285 code: TungsteniteCloseCode::from(code), 286 reason: reason.into(), 287 })), 288 } 289} 290 291fn message_from_tungstenite(message: TungsteniteMessage) -> Result<WsMessage, NetworkError> { 292 match message { 293 TungsteniteMessage::Text(t) => Ok(WsMessage::Text(t.to_string())), 294 TungsteniteMessage::Binary(b) => Ok(WsMessage::Binary(Bytes::copy_from_slice(&b))), 295 TungsteniteMessage::Ping(b) => Ok(WsMessage::Ping(Bytes::copy_from_slice(&b))), 296 TungsteniteMessage::Pong(b) => Ok(WsMessage::Pong(Bytes::copy_from_slice(&b))), 297 TungsteniteMessage::Close(close) => { 298 let (code, reason) = close 299 .map(|c| (u16::from(c.code), c.reason.to_string())) 300 .unwrap_or((1000, String::new())); 301 Ok(WsMessage::Close { code, reason }) 302 } 303 TungsteniteMessage::Frame(_) => Err(NetworkError::Protocol( 304 "tungstenite raw frame surfaced unexpectedly".to_owned(), 305 )), 306 } 307}