I have a warp microservice that in every request, creates a TCP connection and writes/reads from it, the working code looks something like this:
#[derive(Clone, Debug)]
pub struct Redis {
pub host: String,
pub user: Option<String>,
pub pass: Option<String>,
pub v46: bool,
pub port: u16,
pub tls: native_tls::TlsConnector,
}
#[tokio::main]
async fn main() {
let redis: options::Redis = match options::new() {
Ok(o) => o,
Err(e) => {
eprintln!("{}", e);
process::exit(1);
}
};
let args = warp::any().map(move || redis.clone());
let state_route = warp::any()
.and(args)
.and_then(state_handler)
.recover(handle_rejection);
warp::serve(state_route).run((addr, port)).await;
}
And in the state_route per request I do the TCP connection:
async fn state_handler(redis: options::Redis) -> Result<impl warp::Reply, warp::Rejection> {
let conn = timeout(Duration::from_secs(3), TcpStream::connect(&redis.host))
.await
.map_err(|e| warp::reject::custom(RequestTimeout(e.to_string())))?
.map_err(|e| warp::reject::custom(ServiceUnavailable(e.to_string())))?;
let stream = TlsConnector::from(redis.tls.clone())
.connect(&redis.host, conn)
.await
.map_err(|e| warp::reject::custom(ServiceUnavailable(e.to_string())))?;
let mut buf = BufStream::new(stream);
...
}
This works as expected, but I would like to connect on main and share the connection within the handler to prevent multiple connections on every request.
I am trying with Arc something like this:
async fn try_main() -> Result<()> {
let redis: options::Redis = options::new()?;
// TCP connect
let conn = get_connection(&redis.host).await?;
// TLS
let stream = TlsConnector::from(redis.tls.clone())
.connect(&redis.host, conn)
.await?;
let client = Arc::new(stream);
let args = warp::any().map(move || client.clone());
let state_route = warp::any()
.and(args)
.and_then(state_handler)
.recover(handle_rejection);
warp::serve(state_route).run((addr, port)).await;
Ok(())
}
But I don't know how to share the connection to the handler, this is what I am trying:
async fn state_handler(stream: Arc<TlsStream<TcpStream>>) -> Result<impl warp::Reply, warp::Rejection> {
let data = Arc::clone(&stream);
println!("{:#?}", data); // TlsStream... as expected
let mut buf = BufStream::new(data);
...
The error I am getting is when trying to use BufStream, I get:
trait bound `Arc<tokio_native_tls::TlsStream<tokio::net::TcpStream>>: AsyncRead` is not satisfied
|
89 | let mut buf = BufStream::new(data);
| ^^^^ the trait `AsyncWrite` is not implemented for `Arc<tokio_native_tls::TlsStream<tokio::net::TcpStream>>`
|
= note: required by `BufStream::<RW>::new`
I have been suggested to use channels, something like:
let (tx, rx) = mpsc::unbounded_channel();
But I have no idea how could this be implemented within the warp handler since per request I need to write/read to the socket.
Any ideas about how could this be accomplished?