use crate::app_data::ClientAppData;
use crate::bi_stream::BiStream;
use crate::protocol::keepalive::run_keepalive_client_loop;
use crate::quic::client::ClientConfig;
use super::auth::handle_quic_auth_client_side;
use super::metrics_counter::*;
use super::tcp_forwarder::forward_tcp_to_quic_stream;
use anyhow::{Context, Result};
use std::sync::Arc;
use tokio_util::compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt};
use tracing::{debug, info, instrument, warn};
#[instrument(skip(config, conn))]
pub async fn handle_quic_server_connection(
config: Arc<ClientConfig<ClientAppData>>,
conn: quinn::Connection,
) -> Result<()> {
SERVER_CONNECTIONS_OPENED_TOTAL.inc();
let auth_stream = handle_quic_auth_client_side(Arc::clone(&config), conn.clone())
.await
.context("failed to authenticate against PR QUIC server")?;
tokio::spawn(async move {
if let Err(e) = run_keepalive_client_loop(auth_stream).await {
KEEPALIVE_ERRORS.inc();
warn!("Keepalive loop terminated with error: {}", e);
}
});
while let Ok((send, recv)) = conn.accept_bi().await {
CONNECTIONS_ACCEPTED.inc();
let stream_id = recv.id();
let bi_stream = BiStream::new(recv.compat(), send.compat_write(), stream_id.to_string());
info!(
"Opened QUIC stream for new forwarded connection, id {}",
stream_id
);
let config = Arc::clone(&config);
tokio::spawn(async move {
if let Err(e) = forward_tcp_to_quic_stream(config, bi_stream).await {
TCP_FORWARDING_ERRORS.inc();
warn!("Error handling QUIC stream: {}", e);
}
});
}
debug!("Closed QUIC connection handler");
SERVER_CONNECTIONS_GRACEFULLY_CLOSED_TOTAL.inc();
Ok(())
}