how to share a TlsStream<TcpStream> among warp requests

Viewed 202

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?

0 Answers
Related