Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_text/tests/annealer_corpus/hyper_client.rs

21.7 KiB, 1 run

created by r1870400018:11732, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! HTTP/1 client connections
2
3use std::error::Error as StdError;
4use std::fmt;
5use std::future::Future;
6use std::pin::Pin;
7use std::task::{Context, Poll};
8
9use crate::rt::{Read, Write};
10use bytes::Bytes;
11use futures_core::ready;
12use http::{Request, Response};
13use httparse::ParserConfig;
14
15use super::super::dispatch::{self, TrySendError};
16use crate::body::{Body, Incoming as IncomingBody};
17use crate::proto;
18
19type Dispatcher<T, B> =
20 proto::dispatch::Dispatcher<proto::dispatch::Client<B>, B, T, proto::h1::ClientTransaction>;
21
22/// The sender side of an established connection.
23pub struct SendRequest<B> {
24 dispatch: dispatch::Sender<Request<B>, Response<IncomingBody>>,
25}
26
27/// Deconstructed parts of a `Connection`.
28///
29/// This allows taking apart a `Connection` at a later time, in order to
30/// reclaim the IO object, and additional related pieces.
31#[derive(Debug)]
32#[non_exhaustive]
33pub struct Parts<T> {
34 /// The original IO object used in the handshake.
35 pub io: T,
36 /// A buffer of bytes that have been read but not processed as HTTP.
37 ///
38 /// For instance, if the `Connection` is used for an HTTP upgrade request,
39 /// it is possible the server sent back the first bytes of the new protocol
40 /// along with the response upgrade.
41 ///
42 /// You will want to check for any existing bytes if you plan to continue
43 /// communicating on the IO object.
44 pub read_buf: Bytes,
45}
46
47/// A future that processes all HTTP state for the IO object.
48///
49/// In most cases, this should just be spawned into an executor, so that it
50/// can process incoming and outgoing messages, notice hangups, and the like.
51///
52/// Instances of this type are typically created via the [`handshake`] function
53///
54/// # Drop behavior
55///
56/// Dropping the `Connection` will close the underlying IO resource.
57/// Any in-flight requests that have not received a response will be
58/// interrupted. If graceful shutdown is desired, poll the connection
59/// until it completes instead of dropping.
60#[must_use = "futures do nothing unless polled"]
61pub struct Connection<T, B>
62where
63 T: Read + Write,
64 B: Body + 'static,
65{
66 inner: Dispatcher<T, B>,
67}
68
69impl<T, B> Connection<T, B>
70where
71 T: Read + Write + Unpin,
72 B: Body + 'static,
73 B::Error: Into<Box<dyn StdError + Send + Sync>>,
74{
75 /// Return the inner IO object, and additional information.
76 ///
77 /// Only works for HTTP/1 connections. HTTP/2 connections will panic.
78 pub fn into_parts(self) -> Parts<T> {
79 let (io, read_buf, _) = self.inner.into_inner();
80 Parts { io, read_buf }
81 }
82
83 /// Poll the connection for completion, but without calling `shutdown`
84 /// on the underlying IO.
85 ///
86 /// This is useful to allow running a connection while doing an HTTP
87 /// upgrade. Once the upgrade is completed, the connection would be "done",
88 /// but it is not desired to actually shutdown the IO object. Instead you
89 /// would take it back using `into_parts`.
90 ///
91 /// Use [`poll_fn`](https://docs.rs/futures/0.1.25/futures/future/fn.poll_fn.html)
92 /// and [`try_ready!`](https://docs.rs/futures/0.1.25/futures/macro.try_ready.html)
93 /// to work with this function; or use the `without_shutdown` wrapper.
94 pub fn poll_without_shutdown(&mut self, cx: &mut Context<'_>) -> Poll<crate::Result<()>> {
95 self.inner.poll_without_shutdown(cx)
96 }
97
98 /// Prevent shutdown of the underlying IO object at the end of service the request,
99 /// instead run `into_parts`. This is a convenience wrapper over `poll_without_shutdown`.
100 pub async fn without_shutdown(self) -> crate::Result<Parts<T>> {
101 let mut conn = Some(self);
102 crate::common::future::poll_fn(move |cx| -> Poll<crate::Result<Parts<T>>> {
103 ready!(conn.as_mut().unwrap().poll_without_shutdown(cx))?;
104 Poll::Ready(Ok(conn.take().unwrap().into_parts()))
105 })
106 .await
107 }
108}
109
110/// A builder to configure an HTTP connection.
111///
112/// After setting options, the builder is used to create a handshake future.
113///
114/// **Note**: The default values of options are *not considered stable*. They
115/// are subject to change at any time.
116#[derive(Clone, Debug)]
117pub struct Builder {
118 h09_responses: bool,
119 h1_parser_config: ParserConfig,
120 h1_writev: Option<bool>,
121 h1_title_case_headers: bool,
122 h1_preserve_header_case: bool,
123 h1_max_headers: Option<usize>,
124 #[cfg(feature = "ffi")]
125 h1_preserve_header_order: bool,
126 h1_read_buf_exact_size: Option<usize>,
127 h1_max_buf_size: Option<usize>,
128}
129
130/// Returns a handshake future over some IO.
131///
132/// This is a shortcut for `Builder::new().handshake(io)`.
133/// See [`client::conn`](crate::client::conn) for more.
134pub async fn handshake<T, B>(io: T) -> crate::Result<(SendRequest<B>, Connection<T, B>)>
135where
136 T: Read + Write + Unpin,
137 B: Body + 'static,
138 B::Data: Send,
139 B::Error: Into<Box<dyn StdError + Send + Sync>>,
140{
141 Builder::new().handshake(io).await
142}
143
144// ===== impl SendRequest
145
146impl<B> SendRequest<B> {
147 /// Polls to determine whether this sender can be used yet for a request.
148 ///
149 /// If the associated connection is closed, this returns an Error.
150 pub fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<crate::Result<()>> {
151 self.dispatch.poll_ready(cx)
152 }
153
154 /// Waits until the dispatcher is ready
155 ///
156 /// If the associated connection is closed, this returns an Error.
157 pub async fn ready(&mut self) -> crate::Result<()> {
158 crate::common::future::poll_fn(|cx| self.poll_ready(cx)).await
159 }
160
161 /// Checks if the connection is currently ready to send a request.
162 ///
163 /// # Note
164 ///
165 /// This is mostly a hint. Due to inherent latency of networks, it is
166 /// possible that even after checking this is ready, sending a request
167 /// may still fail because the connection was closed in the meantime.
168 pub fn is_ready(&self) -> bool {
169 self.dispatch.is_ready()
170 }
171
172 /// Checks if the connection side has been closed.
173 pub fn is_closed(&self) -> bool {
174 self.dispatch.is_closed()
175 }
176}
177
178impl<B> SendRequest<B>
179where
180 B: Body + 'static,
181{
182 /// Sends a `Request` on the associated connection.
183 ///
184 /// Returns a future that if successful, yields the `Response`.
185 ///
186 /// `req` must have a `Host` header.
187 ///
188 /// # Uri
189 ///
190 /// The `Uri` of the request is serialized as-is.
191 ///
192 /// - Usually you want origin-form (`/path?query`).
193 /// - For sending to an HTTP proxy, you want to send in absolute-form
194 /// (`https://hyper.rs/guides`).
195 ///
196 /// This is however not enforced or validated and it is up to the user
197 /// of this method to ensure the `Uri` is correct for their intended purpose.
198 pub fn send_request(
199 &mut self,
200 req: Request<B>,
201 ) -> impl Future<Output = crate::Result<Response<IncomingBody>>> {
202 let sent = self.dispatch.send(req);
203
204 async move {
205 match sent {
206 Ok(rx) => match rx.await {
207 Ok(Ok(resp)) => Ok(resp),
208 Ok(Err(err)) => Err(err),
209 // this is definite bug if it happens, but it shouldn't happen!
210 Err(_canceled) => panic!("dispatch dropped without returning error"),
211 },
212 Err(_req) => {
213 debug!("connection was not ready");
214 Err(crate::Error::new_canceled().with("connection was not ready"))
215 }
216 }
217 }
218 }
219
220 /// Sends a `Request` on the associated connection.
221 ///
222 /// Returns a future that if successful, yields the `Response`.
223 ///
224 /// # Error
225 ///
226 /// If there was an error before trying to serialize the request to the
227 /// connection, the message will be returned as part of this error.
228 pub fn try_send_request(
229 &mut self,
230 req: Request<B>,
231 ) -> impl Future<Output = Result<Response<IncomingBody>, TrySendError<Request<B>>>> {
232 let sent = self.dispatch.try_send(req);
233 async move {
234 match sent {
235 Ok(rx) => match rx.await {
236 Ok(Ok(res)) => Ok(res),
237 Ok(Err(err)) => Err(err),
238 // this is definite bug if it happens, but it shouldn't happen!
239 Err(_) => panic!("dispatch dropped without returning error"),
240 },
241 Err(req) => {
242 debug!("connection was not ready");
243 let error = crate::Error::new_canceled().with("connection was not ready");
244 Err(TrySendError {
245 error,
246 message: Some(req),
247 })
248 }
249 }
250 }
251 }
252}
253
254impl<B> fmt::Debug for SendRequest<B> {
255 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
256 f.debug_struct("SendRequest").finish()
257 }
258}
259
260// ===== impl Connection
261
262impl<T, B> Connection<T, B>
263where
264 T: Read + Write + Unpin + Send,
265 B: Body + 'static,
266 B::Error: Into<Box<dyn StdError + Send + Sync>>,
267{
268 /// Enable this connection to support higher-level HTTP upgrades.
269 ///
270 /// See [the `upgrade` module](crate::upgrade) for more.
271 pub fn with_upgrades(self) -> upgrades::UpgradeableConnection<T, B> {
272 upgrades::UpgradeableConnection { inner: Some(self) }
273 }
274}
275
276impl<T, B> fmt::Debug for Connection<T, B>
277where
278 T: Read + Write + fmt::Debug,
279 B: Body + 'static,
280{
281 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
282 f.debug_struct("Connection").finish()
283 }
284}
285
286impl<T, B> Future for Connection<T, B>
287where
288 T: Read + Write + Unpin,
289 B: Body + 'static,
290 B::Data: Send,
291 B::Error: Into<Box<dyn StdError + Send + Sync>>,
292{
293 type Output = crate::Result<()>;
294
295 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
296 match ready!(Pin::new(&mut self.inner).poll(cx))? {
297 proto::Dispatched::Shutdown => Poll::Ready(Ok(())),
298 proto::Dispatched::Upgrade(pending) => {
299 // With no `Send` bound on `I`, we can't try to do
300 // upgrades here. In case a user was trying to use
301 // `upgrade` with this API, send a special
302 // error letting them know about that.
303 pending.manual();
304 Poll::Ready(Ok(()))
305 }
306 }
307 }
308}
309
310// ===== impl Builder
311
312impl Builder {
313 /// Creates a new connection builder.
314 #[inline]
315 pub fn new() -> Builder {
316 Builder {
317 h09_responses: false,
318 h1_writev: None,
319 h1_read_buf_exact_size: None,
320 h1_parser_config: Default::default(),
321 h1_title_case_headers: false,
322 h1_preserve_header_case: false,
323 h1_max_headers: None,
324 #[cfg(feature = "ffi")]
325 h1_preserve_header_order: false,
326 h1_max_buf_size: None,
327 }
328 }
329
330 /// Set whether HTTP/0.9 responses should be tolerated.
331 ///
332 /// Default is false.
333 pub fn http09_responses(&mut self, enabled: bool) -> &mut Builder {
334 self.h09_responses = enabled;
335 self
336 }
337
338 /// Set whether HTTP/1 connections will accept spaces between header names
339 /// and the colon that follow them in responses.
340 ///
341 /// You probably don't need this, here is what [RFC 7230 Section 3.2.4.] has
342 /// to say about it:
343 ///
344 /// > No whitespace is allowed between the header field-name and colon. In
345 /// > the past, differences in the handling of such whitespace have led to
346 /// > security vulnerabilities in request routing and response handling. A
347 /// > server MUST reject any received request message that contains
348 /// > whitespace between a header field-name and colon with a response code
349 /// > of 400 (Bad Request). A proxy MUST remove any such whitespace from a
350 /// > response message before forwarding the message downstream.
351 ///
352 /// Default is false.
353 ///
354 /// [RFC 7230 Section 3.2.4.]: https://tools.ietf.org/html/rfc7230#section-3.2.4
355 pub fn allow_spaces_after_header_name_in_responses(&mut self, enabled: bool) -> &mut Builder {
356 self.h1_parser_config
357 .allow_spaces_after_header_name_in_responses(enabled);
358 self
359 }
360
361 /// Set whether HTTP/1 connections will accept obsolete line folding for
362 /// header values.
363 ///
364 /// Newline codepoints (`\r` and `\n`) will be transformed to spaces when
365 /// parsing.
366 ///
367 /// You probably don't need this, here is what [RFC 7230 Section 3.2.4.] has
368 /// to say about it:
369 ///
370 /// > A server that receives an obs-fold in a request message that is not
371 /// > within a message/http container MUST either reject the message by
372 /// > sending a 400 (Bad Request), preferably with a representation
373 /// > explaining that obsolete line folding is unacceptable, or replace
374 /// > each received obs-fold with one or more SP octets prior to
375 /// > interpreting the field value or forwarding the message downstream.
376 ///
377 /// > A proxy or gateway that receives an obs-fold in a response message
378 /// > that is not within a message/http container MUST either discard the
379 /// > message and replace it with a 502 (Bad Gateway) response, preferably
380 /// > with a representation explaining that unacceptable line folding was
381 /// > received, or replace each received obs-fold with one or more SP
382 /// > octets prior to interpreting the field value or forwarding the
383 /// > message downstream.
384 ///
385 /// > A user agent that receives an obs-fold in a response message that is
386 /// > not within a message/http container MUST replace each received
387 /// > obs-fold with one or more SP octets prior to interpreting the field
388 /// > value.
389 ///
390 /// Default is false.
391 ///
392 /// [RFC 7230 Section 3.2.4.]: https://tools.ietf.org/html/rfc7230#section-3.2.4
393 pub fn allow_obsolete_multiline_headers_in_responses(&mut self, enabled: bool) -> &mut Builder {
394 self.h1_parser_config
395 .allow_obsolete_multiline_headers_in_responses(enabled);
396 self
397 }
398
399 /// Set whether HTTP/1 connections will silently ignored malformed header lines.
400 ///
401 /// If this is enabled and a header line does not start with a valid header
402 /// name, or does not include a colon at all, the line will be silently ignored
403 /// and no error will be reported.
404 ///
405 /// Default is false.
406 pub fn ignore_invalid_headers_in_responses(&mut self, enabled: bool) -> &mut Builder {
407 self.h1_parser_config
408 .ignore_invalid_headers_in_responses(enabled);
409 self
410 }
411
412 /// Set whether HTTP/1 connections should try to use vectored writes,
413 /// or always flatten into a single buffer.
414 ///
415 /// Note that setting this to false may mean more copies of body data,
416 /// but may also improve performance when an IO transport doesn't
417 /// support vectored writes well, such as most TLS implementations.
418 ///
419 /// Setting this to true will force hyper to use queued strategy
420 /// which may eliminate unnecessary cloning on some TLS backends
421 ///
422 /// Default is `auto`. In this mode hyper will try to guess which
423 /// mode to use
424 pub fn writev(&mut self, enabled: bool) -> &mut Builder {
425 self.h1_writev = Some(enabled);
426 self
427 }
428
429 /// Set whether HTTP/1 connections will write header names as title case at
430 /// the socket level.
431 ///
432 /// Default is false.
433 pub fn title_case_headers(&mut self, enabled: bool) -> &mut Builder {
434 self.h1_title_case_headers = enabled;
435 self
436 }
437
438 /// Set whether to support preserving original header cases.
439 ///
440 /// Currently, this will record the original cases received, and store them
441 /// in a private extension on the `Response`. It will also look for and use
442 /// such an extension in any provided `Request`.
443 ///
444 /// Since the relevant extension is still private, there is no way to
445 /// interact with the original cases. The only effect this can have now is
446 /// to forward the cases in a proxy-like fashion.
447 ///
448 /// Default is false.
449 pub fn preserve_header_case(&mut self, enabled: bool) -> &mut Builder {
450 self.h1_preserve_header_case = enabled;
451 self
452 }
453
454 /// Set the maximum number of headers.
455 ///
456 /// When a response is received, the parser will reserve a buffer to store headers for optimal
457 /// performance.
458 ///
459 /// If client receives more headers than the buffer size, the error "message header too large"
460 /// is returned.
461 ///
462 /// Note that headers is allocated on the stack by default, which has higher performance. After
463 /// setting this value, headers will be allocated in heap memory, that is, heap memory
464 /// allocation will occur for each response, and there will be a performance drop of about 5%.
465 ///
466 /// Default is 100.
467 pub fn max_headers(&mut self, val: usize) -> &mut Self {
468 self.h1_max_headers = Some(val);
469 self
470 }
471
472 /// Set whether to support preserving original header order.
473 ///
474 /// Currently, this will record the order in which headers are received, and store this
475 /// ordering in a private extension on the `Response`. It will also look for and use
476 /// such an extension in any provided `Request`.
477 ///
478 /// Default is false.
479 #[cfg(feature = "ffi")]
480 pub fn preserve_header_order(&mut self, enabled: bool) -> &mut Builder {
481 self.h1_preserve_header_order = enabled;
482 self
483 }
484
485 /// Sets the exact size of the read buffer to *always* use.
486 ///
487 /// Note that setting this option unsets the `max_buf_size` option.
488 ///
489 /// Default is an adaptive read buffer.
490 pub fn read_buf_exact_size(&mut self, sz: Option<usize>) -> &mut Builder {
491 self.h1_read_buf_exact_size = sz;
492 self.h1_max_buf_size = None;
493 self
494 }
495
496 /// Set the maximum buffer size for the connection.
497 ///
498 /// Default is ~400kb.
499 ///
500 /// Note that setting this option unsets the `read_exact_buf_size` option.
501 ///
502 /// # Panics
503 ///
504 /// The minimum value allowed is 8192. This method panics if the passed `max` is less than the minimum.
505 pub fn max_buf_size(&mut self, max: usize) -> &mut Self {
506 assert!(
507 max >= proto::h1::MINIMUM_MAX_BUFFER_SIZE,
508 "the max_buf_size cannot be smaller than the minimum that h1 specifies."
509 );
510
511 self.h1_max_buf_size = Some(max);
512 self.h1_read_buf_exact_size = None;
513 self
514 }
515
516 /// Constructs a connection with the configured options and IO.
517 /// See [`client::conn`](crate::client::conn) for more.
518 ///
519 /// Note, if [`Connection`] is not `await`-ed, [`SendRequest`] will
520 /// do nothing.
521 pub fn handshake<T, B>(
522 &self,
523 io: T,
524 ) -> impl Future<Output = crate::Result<(SendRequest<B>, Connection<T, B>)>>
525 where
526 T: Read + Write + Unpin,
527 B: Body + 'static,
528 B::Data: Send,
529 B::Error: Into<Box<dyn StdError + Send + Sync>>,
530 {
531 let opts = self.clone();
532
533 async move {
534 trace!("client handshake HTTP/1");
535
536 let (tx, rx) = dispatch::channel();
537 let mut conn = proto::Conn::new(io);
538 conn.set_h1_parser_config(opts.h1_parser_config);
539 if let Some(writev) = opts.h1_writev {
540 if writev {
541 conn.set_write_strategy_queue();
542 } else {
543 conn.set_write_strategy_flatten();
544 }
545 }
546 if opts.h1_title_case_headers {
547 conn.set_title_case_headers();
548 }
549 if opts.h1_preserve_header_case {
550 conn.set_preserve_header_case();
551 }
552 if let Some(max_headers) = opts.h1_max_headers {
553 conn.set_http1_max_headers(max_headers);
554 }
555 #[cfg(feature = "ffi")]
556 if opts.h1_preserve_header_order {
557 conn.set_preserve_header_order();
558 }
559
560 if opts.h09_responses {
561 conn.set_h09_responses();
562 }
563
564 if let Some(sz) = opts.h1_read_buf_exact_size {
565 conn.set_read_buf_exact_size(sz);
566 }
567 if let Some(max) = opts.h1_max_buf_size {
568 conn.set_max_buf_size(max);
569 }
570 let cd = proto::h1::dispatch::Client::new(rx);
571 let proto = proto::h1::Dispatcher::new(cd, conn);
572
573 Ok((SendRequest { dispatch: tx }, Connection { inner: proto }))
574 }
575 }
576}
577
578mod upgrades {
579 use crate::upgrade::Upgraded;
580
581 use super::*;
582
583 // A future binding a connection with a Service with Upgrade support.
584 //
585 // This type is unnameable outside the crate.
586 #[must_use = "futures do nothing unless polled"]
587 #[allow(missing_debug_implementations)]
588 pub struct UpgradeableConnection<T, B>
589 where
590 T: Read + Write + Unpin + Send + 'static,
591 B: Body + 'static,
592 B::Error: Into<Box<dyn StdError + Send + Sync>>,
593 {
594 pub(super) inner: Option<Connection<T, B>>,
595 }
596
597 impl<I, B> Future for UpgradeableConnection<I, B>
598 where
599 I: Read + Write + Unpin + Send + 'static,
600 B: Body + 'static,
601 B::Data: Send,
602 B::Error: Into<Box<dyn StdError + Send + Sync>>,
603 {
604 type Output = crate::Result<()>;
605
606 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
607 match ready!(Pin::new(&mut self.inner.as_mut().unwrap().inner).poll(cx)) {
608 Ok(proto::Dispatched::Shutdown) => Poll::Ready(Ok(())),
609 Ok(proto::Dispatched::Upgrade(pending)) => {
610 let Parts { io, read_buf } = self.inner.take().unwrap().into_parts();
611 pending.fulfill(Upgraded::new(io, read_buf));
612 Poll::Ready(Ok(()))
613 }
614 Err(e) => Poll::Ready(Err(e)),
615 }
616 }
617 }
618}