use std::error::Error as StdError; use std::future::Future; use std::pin::Pin; use std::task::{Context, Poll}; use std::time::Duration;
use bytes::Bytes; use futures_core::ready; use h2::server::{Connection, Handshake, SendResponse}; use h2::{Reason, RecvStream}; use http::{Method, Request}; use pin_project_lite::pin_project;
// Our defaults are chosen for the "majority" case, which usually are not // resource constrained, and so the spec default of 64kb can be too limiting // for performance. // // At the same time, a server more often has multiple clients connected, and // so is more likely to use more resources than a client would. const DEFAULT_CONN_WINDOW: u32 = 1024 * 1024; // 1mb const DEFAULT_STREAM_WINDOW: u32 = 1024 * 1024; // 1mb const DEFAULT_MAX_FRAME_SIZE: u32 = 1024 * 16; // 16kb const DEFAULT_MAX_SEND_BUF_SIZE: usize = 1024 * 400; // 400kb const DEFAULT_SETTINGS_MAX_HEADER_LIST_SIZE: u32 = 1024 * 16; // 16kb const DEFAULT_MAX_LOCAL_ERROR_RESET_STREAMS: usize = 1024;
let bdp = if config.adaptive_window {
Some(config.initial_stream_window_size)
} else {
None
};
let ping_config = ping::Config {
bdp_initial_window: bdp,
keep_alive_interval: config.keep_alive_interval,
keep_alive_timeout: config.keep_alive_timeout, // If keep-alive is enabled for servers, always enabled while // idle, so it can more aggressively close dead connections.
keep_alive_while_idle: true,
};
impl<F, B, Ex, E> H2Stream<F, B, Ex> where
F: Future<Output = Result<Response<B>, E>>,
B: Body,
B::Data: 'static,
B::Error: Into<Box<dyn StdError + Send + Sync>>,
Ex: Http2UpgradedExec<B::Data>,
E: Into<Box<dyn StdError + Send + Sync>>,
{ fn poll2(mutself: Pin<&mutSelf>, cx: &mut Context<'_>) -> Poll<crate::Result<()>> { letmut me = self.as_mut().project(); loop { let next = match me.state.as_mut().project() {
H2StreamStateProj::Service {
fut: h,
connect_parts,
} => { let res = match h.poll(cx) {
Poll::Ready(Ok(r)) => r,
Poll::Pending => { // Response is not yet ready, so we want to check if the client has sent a // RST_STREAM frame which would cancel the current request. iflet Poll::Ready(reason) =
me.reply.poll_reset(cx).map_err(crate::Error::new_h2)?
{
debug!("stream received RST_STREAM: {:?}", reason); return Poll::Ready(Err(crate::Error::new_h2(reason.into())));
} return Poll::Pending;
}
Poll::Ready(Err(e)) => { let err = crate::Error::new_user_service(e);
warn!("http2 service errored: {}", err);
me.reply.send_reset(err.h2_reason()); return Poll::Ready(Err(err));
}
};
let (head, body) = res.into_parts(); letmut res = ::http::Response::from_parts(head, ()); super::strip_connection_headers(res.headers_mut(), false);
// set Date header if it isn't already set if instructed if *me.date_header {
res.headers_mut()
.entry(::http::header::DATE)
.or_insert_with(date::update_and_header_value);
}
iflet Some(connect_parts) = connect_parts.take() { if res.status().is_success() { if headers::content_length_parse_all(res.headers())
.map_or(false, |len| len != 0)
{
warn!("h2 successful response to CONNECT request with body not supported");
me.reply.send_reset(h2::Reason::INTERNAL_ERROR); return Poll::Ready(Err(crate::Error::new_user_header()));
} if res
.headers_mut()
.remove(::http::header::CONTENT_LENGTH)
.is_some()
{
warn!("successful response to CONNECT request disallows content-length header");
} let send_stream = reply!(me, res, false); let (h2_up, up_task) = super::upgrade::pair(
send_stream,
connect_parts.recv_stream,
connect_parts.ping,
);
connect_parts
.pending
.fulfill(Upgraded::new(h2_up, Bytes::new())); self.exec.execute_upgrade(up_task); return Poll::Ready(Ok(()));
}
}
if !body.is_end_stream() { // automatically set Content-Length from body... iflet Some(len) = body.size_hint().exact() {
headers::set_content_length_if_missing(res.headers_mut(), len);
}
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.