fix bug and somewhat reusable streams
This commit is contained in:
parent
39493beb53
commit
aa95ab9c16
3 changed files with 19 additions and 12 deletions
|
|
@ -10,7 +10,6 @@ where
|
||||||
T: Connect + Clone + Sync + Send + 'static,
|
T: Connect + Clone + Sync + Send + 'static,
|
||||||
{
|
{
|
||||||
pub async fn client_serve_h3(self, conn: quinn::Connecting) -> Result<()> {
|
pub async fn client_serve_h3(self, conn: quinn::Connecting) -> Result<()> {
|
||||||
// TODO: client数の管理
|
|
||||||
let client_addr = conn.remote_address();
|
let client_addr = conn.remote_address();
|
||||||
|
|
||||||
match conn.await {
|
match conn.await {
|
||||||
|
|
@ -38,10 +37,10 @@ where
|
||||||
|
|
||||||
// TODO: Work around for timeout...
|
// TODO: Work around for timeout...
|
||||||
// while let Some((req, stream)) = h3_conn
|
// while let Some((req, stream)) = h3_conn
|
||||||
// if let Some((req, stream)) =
|
// .accept()
|
||||||
// .await
|
// .await
|
||||||
// .map_err(|e| anyhow!("HTTP/3 accept failed: {}", e))?
|
// .map_err(|e| anyhow!("HTTP/3 accept failed: {}", e))?
|
||||||
if let Some((req, stream)) = match tokio::time::timeout(
|
while let Some((req, stream)) = match tokio::time::timeout(
|
||||||
tokio::time::Duration::from_millis(H3_CONN_TIMEOUT_MILLIS),
|
tokio::time::Duration::from_millis(H3_CONN_TIMEOUT_MILLIS),
|
||||||
h3_conn.accept(),
|
h3_conn.accept(),
|
||||||
)
|
)
|
||||||
|
|
@ -49,7 +48,7 @@ where
|
||||||
{
|
{
|
||||||
Ok(r) => r.map_err(|e| anyhow!("HTTP/3 accept failed: {}", e))?,
|
Ok(r) => r.map_err(|e| anyhow!("HTTP/3 accept failed: {}", e))?,
|
||||||
Err(_) => {
|
Err(_) => {
|
||||||
warn!("No incoming stream after connection establishment");
|
warn!("No incoming stream after connection establishment / previous use");
|
||||||
h3_conn.shutdown(0).await?;
|
h3_conn.shutdown(0).await?;
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
|
|
@ -66,11 +65,11 @@ where
|
||||||
if let Err(e) = self_inner.handle_request_h3(req, stream, client_addr).await {
|
if let Err(e) = self_inner.handle_request_h3(req, stream, client_addr).await {
|
||||||
error!("HTTP/3 request failed: {}", e);
|
error!("HTTP/3 request failed: {}", e);
|
||||||
}
|
}
|
||||||
// TODO: Work around for timeout
|
// // TODO: Work around for timeout
|
||||||
if let Err(e) = h3_conn.shutdown(0).await {
|
// if let Err(e) = h3_conn.shutdown(0).await {
|
||||||
error!("HTTP/3 connection shutdown failed: {}", e);
|
// error!("HTTP/3 connection shutdown failed: {}", e);
|
||||||
}
|
// }
|
||||||
debug!("HTTP/3 connection shutdown (currently shutdown each time as work around for timeout)");
|
// debug!("HTTP/3 connection shutdown (currently shutdown each time as work around for timeout)");
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -115,7 +115,7 @@ where
|
||||||
res_backend.headers_mut().insert(
|
res_backend.headers_mut().insert(
|
||||||
hyper::header::ALT_SVC,
|
hyper::header::ALT_SVC,
|
||||||
format!(
|
format!(
|
||||||
"h3=\":{}\"; ma={}, h3-29\":{}\"; ma={}",
|
"h3=\":{}\"; ma={}, h3-29=\":{}\"; ma={}",
|
||||||
port, H3_ALT_SVC_MAX_AGE, port, H3_ALT_SVC_MAX_AGE
|
port, H3_ALT_SVC_MAX_AGE, port, H3_ALT_SVC_MAX_AGE
|
||||||
)
|
)
|
||||||
.parse()
|
.parse()
|
||||||
|
|
|
||||||
|
|
@ -144,11 +144,19 @@ where
|
||||||
let peekable_incoming = std::pin::Pin::new(&mut p);
|
let peekable_incoming = std::pin::Pin::new(&mut p);
|
||||||
if let Some(conn) = peekable_incoming.get_mut().next().await {
|
if let Some(conn) = peekable_incoming.get_mut().next().await {
|
||||||
if success {
|
if success {
|
||||||
|
// TODO: client数の管理
|
||||||
|
let clients_count = self.globals.clients_count.clone();
|
||||||
|
if clients_count.increment() > self.globals.max_clients {
|
||||||
|
clients_count.decrement();
|
||||||
|
continue;
|
||||||
|
}
|
||||||
let fut = self.clone().client_serve_h3(conn);
|
let fut = self.clone().client_serve_h3(conn);
|
||||||
self.globals.runtime_handle.spawn(async {
|
self.globals.runtime_handle.spawn(async move {
|
||||||
if let Err(e) = fut.await {
|
if let Err(e) = fut.await {
|
||||||
warn!("QUIC or HTTP/3 connection failed: {}", e)
|
warn!("QUIC or HTTP/3 connection failed: {}", e)
|
||||||
}
|
}
|
||||||
|
clients_count.decrement();
|
||||||
|
debug!("Client #: {}", clients_count.current());
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue