1#
2# Copyright 2014 Facebook
3#
4# Licensed under the Apache License, Version 2.0 (the "License"); you may
5# not use this file except in compliance with the License. You may obtain
6# a copy of the License at
7#
8# http://www.apache.org/licenses/LICENSE-2.0
9#
10# Unless required by applicable law or agreed to in writing, software
11# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
12# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
13# License for the specific language governing permissions and limitations
14# under the License.
15
16"""Client and server implementations of HTTP/1.x.
17
18.. versionadded:: 4.0
19"""
20
21import asyncio
22import logging
23import re
24import types
25
26from tornado.concurrent import (
27 Future,
28 future_add_done_callback,
29 future_set_result_unless_cancelled,
30)
31from tornado.escape import native_str, utf8
32from tornado import gen
33from tornado import httputil
34from tornado import iostream
35from tornado.log import gen_log, app_log
36from tornado.util import GzipDecompressor
37
38
39from typing import cast, Optional, Type, Awaitable, Callable, Union, Tuple
40
41CR_OR_LF_RE = re.compile(b"\r|\n")
42
43
44# The maximum number of informational (1xx) responses to accept before the
45# real response. Each one is processed with a recursive call to
46# _read_message, so an unbounded number of them would exhaust the stack.
47# There is no legitimate use for more than a handful.
48_MAX_1XX_RESPONSES = 10
49
50
51class _QuietException(Exception):
52 def __init__(self) -> None:
53 pass
54
55
56class _ExceptionLoggingContext:
57 """Used with the ``with`` statement when calling delegate methods to
58 log any exceptions with the given logger. Any exceptions caught are
59 converted to _QuietException
60 """
61
62 def __init__(self, logger: logging.Logger) -> None:
63 self.logger = logger
64
65 def __enter__(self) -> None:
66 pass
67
68 def __exit__(
69 self,
70 typ: "Optional[Type[BaseException]]",
71 value: Optional[BaseException],
72 tb: types.TracebackType,
73 ) -> None:
74 if value is not None:
75 assert typ is not None
76 # Let HTTPInputError pass through to higher-level handler
77 if isinstance(value, httputil.HTTPInputError):
78 return None
79 self.logger.error("Uncaught exception", exc_info=(typ, value, tb))
80 raise _QuietException
81
82
83class HTTP1ConnectionParameters:
84 """Parameters for `.HTTP1Connection` and `.HTTP1ServerConnection`."""
85
86 def __init__(
87 self,
88 no_keep_alive: bool = False,
89 chunk_size: Optional[int] = None,
90 max_header_size: Optional[int] = None,
91 header_timeout: Optional[float] = None,
92 max_body_size: Optional[int] = None,
93 body_timeout: Optional[float] = None,
94 decompress: bool = False,
95 ) -> None:
96 """
97 :arg bool no_keep_alive: If true, always close the connection after
98 one request.
99 :arg int chunk_size: how much data to read into memory at once
100 :arg int max_header_size: maximum amount of data for HTTP headers
101 :arg float header_timeout: how long to wait for all headers (seconds)
102 :arg int max_body_size: maximum amount of data for body
103 :arg float body_timeout: how long to wait while reading body (seconds)
104 :arg bool decompress: if true, decode incoming
105 ``Content-Encoding: gzip``
106 """
107 self.no_keep_alive = no_keep_alive
108 self.chunk_size = chunk_size or 65536
109 self.max_header_size = max_header_size or 65536
110 self.header_timeout = header_timeout
111 self.max_body_size = max_body_size
112 self.body_timeout = body_timeout
113 self.decompress = decompress
114
115
116class HTTP1Connection(httputil.HTTPConnection):
117 """Implements the HTTP/1.x protocol.
118
119 This class can be on its own for clients, or via `HTTP1ServerConnection`
120 for servers.
121 """
122
123 def __init__(
124 self,
125 stream: iostream.IOStream,
126 is_client: bool,
127 params: Optional[HTTP1ConnectionParameters] = None,
128 context: Optional[object] = None,
129 ) -> None:
130 """
131 :arg stream: an `.IOStream`
132 :arg bool is_client: client or server
133 :arg params: a `.HTTP1ConnectionParameters` instance or ``None``
134 :arg context: an opaque application-defined object that can be accessed
135 as ``connection.context``.
136 """
137 self.is_client = is_client
138 self.stream = stream
139 if params is None:
140 params = HTTP1ConnectionParameters()
141 self.params = params
142 self.context = context
143 self.no_keep_alive = params.no_keep_alive
144 # The body limits can be altered by the delegate, so save them
145 # here instead of just referencing self.params later.
146 self._max_body_size = (
147 self.params.max_body_size
148 if self.params.max_body_size is not None
149 else self.stream.max_buffer_size
150 )
151 self._body_timeout = self.params.body_timeout
152 # _write_finished is set to True when finish() has been called,
153 # i.e. there will be no more data sent. Data may still be in the
154 # stream's write buffer.
155 self._write_finished = False
156 # True when we have read the entire incoming body.
157 self._read_finished = False
158 # _finish_future resolves when all data has been written and flushed
159 # to the IOStream.
160 self._finish_future = Future() # type: Future[None]
161 # If true, the connection should be closed after this request
162 # (after the response has been written in the server side,
163 # and after it has been read in the client)
164 self._disconnect_on_finish = False
165 self._clear_callbacks()
166 # Save the start lines after we read or write them; they
167 # affect later processing (e.g. 304 responses and HEAD methods
168 # have content-length but no bodies)
169 self._request_start_line = None # type: Optional[httputil.RequestStartLine]
170 self._response_start_line = None # type: Optional[httputil.ResponseStartLine]
171 self._request_headers = None # type: Optional[httputil.HTTPHeaders]
172 # True if we are writing output with chunked encoding.
173 self._chunking_output = False
174 # While reading a body with a content-length, this is the
175 # amount left to read.
176 self._expected_content_remaining = None # type: Optional[int]
177 # A Future for our outgoing writes, returned by IOStream.write.
178 self._pending_write = None # type: Optional[Future[None]]
179
180 def read_response(self, delegate: httputil.HTTPMessageDelegate) -> Awaitable[bool]:
181 """Read a single HTTP response.
182
183 Typical client-mode usage is to write a request using `write_headers`,
184 `write`, and `finish`, and then call ``read_response``.
185
186 :arg delegate: a `.HTTPMessageDelegate`
187
188 Returns a `.Future` that resolves to a bool after the full response has
189 been read. The result is true if the stream is still open.
190 """
191 if self.params.decompress:
192 delegate = _GzipMessageDelegate(
193 delegate, self.params.chunk_size, self._max_body_size
194 )
195 return self._read_message(delegate)
196
197 async def _read_message(
198 self, delegate: httputil.HTTPMessageDelegate, num_1xx: int = 0
199 ) -> bool:
200 need_delegate_close = False
201 try:
202 header_future = self.stream.read_until_regex(
203 b"\r?\n\r?\n", max_bytes=self.params.max_header_size
204 )
205 if self.params.header_timeout is None:
206 header_data = await header_future
207 else:
208 try:
209 header_data = await gen.with_timeout(
210 self.stream.io_loop.time() + self.params.header_timeout,
211 header_future,
212 quiet_exceptions=iostream.StreamClosedError,
213 )
214 except gen.TimeoutError:
215 self.close()
216 return False
217 start_line_str, headers = self._parse_headers(header_data)
218 if self.is_client:
219 resp_start_line = httputil.parse_response_start_line(start_line_str)
220 self._response_start_line = resp_start_line
221 start_line = (
222 resp_start_line
223 ) # type: Union[httputil.RequestStartLine, httputil.ResponseStartLine]
224 # TODO: this will need to change to support client-side keepalive
225 self._disconnect_on_finish = False
226 else:
227 req_start_line = httputil.parse_request_start_line(start_line_str)
228 self._request_start_line = req_start_line
229 self._request_headers = headers
230 start_line = req_start_line
231 self._disconnect_on_finish = not self._can_keep_alive(
232 req_start_line, headers
233 )
234 need_delegate_close = True
235 with _ExceptionLoggingContext(app_log):
236 header_recv_future = delegate.headers_received(start_line, headers)
237 if header_recv_future is not None:
238 await header_recv_future
239 if self.stream is None:
240 # We've been detached.
241 need_delegate_close = False
242 return False
243 skip_body = False
244 if self.is_client:
245 assert isinstance(start_line, httputil.ResponseStartLine)
246 if (
247 self._request_start_line is not None
248 and self._request_start_line.method == "HEAD"
249 ):
250 skip_body = True
251 code = start_line.code
252 if code == 304:
253 # 304 responses may include the content-length header
254 # but do not actually have a body.
255 # http://tools.ietf.org/html/rfc7230#section-3.3
256 skip_body = True
257 if 100 <= code < 200:
258 # 1xx responses should never indicate the presence of
259 # a body.
260 if "Content-Length" in headers or "Transfer-Encoding" in headers:
261 raise httputil.HTTPInputError(
262 "Response code %d cannot have body" % code
263 )
264 if num_1xx >= _MAX_1XX_RESPONSES:
265 raise httputil.HTTPInputError("Too many 1xx responses")
266 # TODO: client delegates will get headers_received twice
267 # in the case of a 100-continue. Document or change?
268 #
269 # The recursive call reads the real response and owns
270 # the delegate from here on, so there is nothing left
271 # for this frame to do. Clear need_delegate_close so
272 # that the finally block does not call
273 # on_connection_close() on an already-finished
274 # delegate.
275 need_delegate_close = False
276 return await self._read_message(delegate, num_1xx + 1)
277 else:
278 if headers.get("Expect") == "100-continue" and not self._write_finished:
279 self.stream.write(b"HTTP/1.1 100 (Continue)\r\n\r\n")
280 if not skip_body:
281 body_future = self._read_body(
282 resp_start_line.code if self.is_client else 0, headers, delegate
283 )
284 if body_future is not None:
285 if self._body_timeout is None:
286 await body_future
287 else:
288 try:
289 await gen.with_timeout(
290 self.stream.io_loop.time() + self._body_timeout,
291 body_future,
292 quiet_exceptions=iostream.StreamClosedError,
293 )
294 except gen.TimeoutError:
295 gen_log.info("Timeout reading body from %s", self.context)
296 self.stream.close()
297 return False
298 self._read_finished = True
299 if not self._write_finished or self.is_client:
300 need_delegate_close = False
301 with _ExceptionLoggingContext(app_log):
302 delegate.finish()
303 # If we're waiting for the application to produce an asynchronous
304 # response, and we're not detached, register a close callback
305 # on the stream (we didn't need one while we were reading)
306 if (
307 not self._finish_future.done()
308 and self.stream is not None
309 and not self.stream.closed()
310 ):
311 self.stream.set_close_callback(self._on_connection_close)
312 await self._finish_future
313 if self.is_client and self._disconnect_on_finish:
314 self.close()
315 if self.stream is None:
316 return False
317 except httputil.HTTPInputError as e:
318 gen_log.info("Malformed HTTP message from %s: %s", self.context, e)
319 if not self.is_client:
320 await self.stream.write(b"HTTP/1.1 400 Bad Request\r\n\r\n")
321 self.close()
322 return False
323 finally:
324 if need_delegate_close:
325 with _ExceptionLoggingContext(app_log):
326 delegate.on_connection_close()
327 header_future = None # type: ignore
328 self._clear_callbacks()
329 return True
330
331 def _clear_callbacks(self) -> None:
332 """Clears the callback attributes.
333
334 This allows the request handler to be garbage collected more
335 quickly in CPython by breaking up reference cycles.
336 """
337 self._write_callback = None
338 self._write_future = None # type: Optional[Future[None]]
339 self._close_callback = None # type: Optional[Callable[[], None]]
340 if self.stream is not None:
341 self.stream.set_close_callback(None)
342
343 def set_close_callback(self, callback: Optional[Callable[[], None]]) -> None:
344 """Sets a callback that will be run when the connection is closed.
345
346 Note that this callback is slightly different from
347 `.HTTPMessageDelegate.on_connection_close`: The
348 `.HTTPMessageDelegate` method is called when the connection is
349 closed while receiving a message. This callback is used when
350 there is not an active delegate (for example, on the server
351 side this callback is used if the client closes the connection
352 after sending its request but before receiving all the
353 response.
354 """
355 self._close_callback = callback
356
357 def _on_connection_close(self) -> None:
358 # Note that this callback is only registered on the IOStream
359 # when we have finished reading the request and are waiting for
360 # the application to produce its response.
361 if self._close_callback is not None:
362 callback = self._close_callback
363 self._close_callback = None
364 callback()
365 if not self._finish_future.done():
366 future_set_result_unless_cancelled(self._finish_future, None)
367 self._clear_callbacks()
368
369 def close(self) -> None:
370 if self.stream is not None:
371 self.stream.close()
372 self._clear_callbacks()
373 if not self._finish_future.done():
374 future_set_result_unless_cancelled(self._finish_future, None)
375
376 def detach(self) -> iostream.IOStream:
377 """Take control of the underlying stream.
378
379 Returns the underlying `.IOStream` object and stops all further
380 HTTP processing. May only be called during
381 `.HTTPMessageDelegate.headers_received`. Intended for implementing
382 protocols like websockets that tunnel over an HTTP handshake.
383 """
384 self._clear_callbacks()
385 stream = self.stream
386 self.stream = None # type: ignore
387 if not self._finish_future.done():
388 future_set_result_unless_cancelled(self._finish_future, None)
389 return stream
390
391 def set_body_timeout(self, timeout: float) -> None:
392 """Sets the body timeout for a single request.
393
394 Overrides the value from `.HTTP1ConnectionParameters`.
395 """
396 self._body_timeout = timeout
397
398 def set_max_body_size(self, max_body_size: int) -> None:
399 """Sets the body size limit for a single request.
400
401 Overrides the value from `.HTTP1ConnectionParameters`.
402 """
403 self._max_body_size = max_body_size
404
405 def write_headers(
406 self,
407 start_line: Union[httputil.RequestStartLine, httputil.ResponseStartLine],
408 headers: httputil.HTTPHeaders,
409 chunk: Optional[bytes] = None,
410 ) -> "Future[None]":
411 """Implements `.HTTPConnection.write_headers`."""
412 lines = []
413 if self.is_client:
414 assert isinstance(start_line, httputil.RequestStartLine)
415 self._request_start_line = start_line
416 lines.append(utf8(f"{start_line[0]} {start_line[1]} HTTP/1.1"))
417 # Client requests with a non-empty body must have either a
418 # Content-Length or a Transfer-Encoding. If Content-Length is not
419 # present we'll add our Transfer-Encoding below.
420 self._chunking_output = (
421 start_line.method in ("POST", "PUT", "PATCH")
422 and "Content-Length" not in headers
423 )
424 else:
425 assert isinstance(start_line, httputil.ResponseStartLine)
426 assert self._request_start_line is not None
427 assert self._request_headers is not None
428 self._response_start_line = start_line
429 lines.append(utf8("HTTP/1.1 %d %s" % (start_line[1], start_line[2])))
430 self._chunking_output = (
431 # TODO: should this use
432 # self._request_start_line.version or
433 # start_line.version?
434 self._request_start_line.version == "HTTP/1.1"
435 # Omit payload header field for HEAD request.
436 and self._request_start_line.method != "HEAD"
437 # 1xx, 204 and 304 responses have no body (not even a zero-length
438 # body), and so should not have either Content-Length or
439 # Transfer-Encoding headers.
440 and start_line.code not in (204, 304)
441 and (start_line.code < 100 or start_line.code >= 200)
442 # No need to chunk the output if a Content-Length is specified.
443 and "Content-Length" not in headers
444 )
445 # If connection to a 1.1 client will be closed, inform client
446 if (
447 self._request_start_line.version == "HTTP/1.1"
448 and self._disconnect_on_finish
449 ):
450 headers["Connection"] = "close"
451 # If a 1.0 client asked for keep-alive, add the header.
452 if (
453 self._request_start_line.version == "HTTP/1.0"
454 and self._request_headers.get("Connection", "").lower() == "keep-alive"
455 ):
456 headers["Connection"] = "Keep-Alive"
457 if self._chunking_output:
458 headers["Transfer-Encoding"] = "chunked"
459 if not self.is_client and (
460 self._request_start_line.method == "HEAD"
461 or cast(httputil.ResponseStartLine, start_line).code == 304
462 ):
463 self._expected_content_remaining = 0
464 elif "Content-Length" in headers:
465 self._expected_content_remaining = parse_int(headers["Content-Length"])
466 else:
467 self._expected_content_remaining = None
468 # TODO: headers are supposed to be of type str, but we still have some
469 # cases that let bytes slip through. Remove these native_str calls when those
470 # are fixed.
471 header_lines = (
472 native_str(n) + ": " + native_str(v) for n, v in headers.get_all()
473 )
474 lines.extend(line.encode("latin1") for line in header_lines)
475 for line in lines:
476 if CR_OR_LF_RE.search(line):
477 raise ValueError("Illegal characters (CR or LF) in header: %r" % line)
478 future = None
479 if self.stream.closed():
480 future = self._write_future = Future()
481 future.set_exception(iostream.StreamClosedError())
482 future.exception()
483 else:
484 future = self._write_future = Future()
485 data = b"\r\n".join(lines) + b"\r\n\r\n"
486 if chunk:
487 data += self._format_chunk(chunk)
488 self._pending_write = self.stream.write(data)
489 future_add_done_callback(self._pending_write, self._on_write_complete)
490 return future
491
492 def _format_chunk(self, chunk: bytes) -> bytes:
493 if self._expected_content_remaining is not None:
494 self._expected_content_remaining -= len(chunk)
495 if self._expected_content_remaining < 0:
496 # Close the stream now to stop further framing errors.
497 self.stream.close()
498 raise httputil.HTTPOutputError(
499 "Tried to write more data than Content-Length"
500 )
501 if self._chunking_output and chunk:
502 # Don't write out empty chunks because that means END-OF-STREAM
503 # with chunked encoding
504 return utf8("%x" % len(chunk)) + b"\r\n" + chunk + b"\r\n"
505 else:
506 return chunk
507
508 def write(self, chunk: bytes) -> "Future[None]":
509 """Implements `.HTTPConnection.write`.
510
511 For backwards compatibility it is allowed but deprecated to
512 skip `write_headers` and instead call `write()` with a
513 pre-encoded header block.
514 """
515 future = None
516 if self.stream.closed():
517 future = self._write_future = Future()
518 self._write_future.set_exception(iostream.StreamClosedError())
519 self._write_future.exception()
520 else:
521 future = self._write_future = Future()
522 self._pending_write = self.stream.write(self._format_chunk(chunk))
523 future_add_done_callback(self._pending_write, self._on_write_complete)
524 return future
525
526 def finish(self) -> None:
527 """Implements `.HTTPConnection.finish`."""
528 if (
529 self._expected_content_remaining is not None
530 and self._expected_content_remaining != 0
531 and not self.stream.closed()
532 ):
533 self.stream.close()
534 raise httputil.HTTPOutputError(
535 "Tried to write %d bytes less than Content-Length"
536 % self._expected_content_remaining
537 )
538 if self._chunking_output:
539 if not self.stream.closed():
540 self._pending_write = self.stream.write(b"0\r\n\r\n")
541 self._pending_write.add_done_callback(self._on_write_complete)
542 self._write_finished = True
543 # If the app finished the request while we're still reading,
544 # divert any remaining data away from the delegate and
545 # close the connection when we're done sending our response.
546 # Closing the connection is the only way to avoid reading the
547 # whole input body.
548 if not self._read_finished:
549 self._disconnect_on_finish = True
550 # No more data is coming, so instruct TCP to send any remaining
551 # data immediately instead of waiting for a full packet or ack.
552 self.stream.set_nodelay(True)
553 if self._pending_write is None:
554 self._finish_request(None)
555 else:
556 future_add_done_callback(self._pending_write, self._finish_request)
557
558 def _on_write_complete(self, future: "Future[None]") -> None:
559 exc = future.exception()
560 if exc is not None and not isinstance(exc, iostream.StreamClosedError):
561 future.result()
562 if self._write_callback is not None:
563 callback = self._write_callback
564 self._write_callback = None
565 self.stream.io_loop.add_callback(callback)
566 if self._write_future is not None:
567 future = self._write_future
568 self._write_future = None
569 future_set_result_unless_cancelled(future, None)
570
571 def _can_keep_alive(
572 self, start_line: httputil.RequestStartLine, headers: httputil.HTTPHeaders
573 ) -> bool:
574 if self.params.no_keep_alive:
575 return False
576 connection_header = headers.get("Connection")
577 if connection_header is not None:
578 connection_header = connection_header.lower()
579 if start_line.version == "HTTP/1.1":
580 return connection_header != "close"
581 elif (
582 "Content-Length" in headers
583 or is_transfer_encoding_chunked(headers)
584 or getattr(start_line, "method", None) in ("HEAD", "GET")
585 ):
586 # start_line may be a request or response start line; only
587 # the former has a method attribute.
588 return connection_header == "keep-alive"
589 return False
590
591 def _finish_request(self, future: "Optional[Future[None]]") -> None:
592 self._clear_callbacks()
593 if not self.is_client and self._disconnect_on_finish:
594 self.close()
595 return
596 # Turn Nagle's algorithm back on, leaving the stream in its
597 # default state for the next request.
598 self.stream.set_nodelay(False)
599 if not self._finish_future.done():
600 future_set_result_unless_cancelled(self._finish_future, None)
601
602 def _parse_headers(self, data: bytes) -> Tuple[str, httputil.HTTPHeaders]:
603 # The lstrip removes newlines that some implementations sometimes
604 # insert between messages of a reused connection. Per RFC 7230,
605 # we SHOULD ignore at least one empty line before the request.
606 # http://tools.ietf.org/html/rfc7230#section-3.5
607 data_str = native_str(data.decode("latin1")).lstrip("\r\n")
608 # RFC 7230 section allows for both CRLF and bare LF.
609 eol = data_str.find("\n")
610 start_line = data_str[:eol].rstrip("\r")
611 headers = httputil.HTTPHeaders.parse(data_str[eol:])
612 return start_line, headers
613
614 def _read_body(
615 self,
616 code: int,
617 headers: httputil.HTTPHeaders,
618 delegate: httputil.HTTPMessageDelegate,
619 ) -> Optional[Awaitable[None]]:
620 if "Content-Length" in headers:
621 if "," in headers["Content-Length"]:
622 # Proxies sometimes cause Content-Length headers to get
623 # duplicated. If all the values are identical then we can
624 # use them but if they differ it's an error.
625 pieces = re.split(r",\s*", headers["Content-Length"])
626 if any(i != pieces[0] for i in pieces):
627 raise httputil.HTTPInputError(
628 "Multiple unequal Content-Lengths: %r"
629 % headers["Content-Length"]
630 )
631 headers["Content-Length"] = pieces[0]
632
633 try:
634 content_length: Optional[int] = parse_int(headers["Content-Length"])
635 except ValueError:
636 # Handles non-integer Content-Length value.
637 raise httputil.HTTPInputError(
638 "Only integer Content-Length is allowed: %s"
639 % headers["Content-Length"]
640 )
641
642 if cast(int, content_length) > self._max_body_size:
643 raise httputil.HTTPInputError("Content-Length too long")
644 else:
645 content_length = None
646
647 is_chunked = is_transfer_encoding_chunked(headers)
648
649 if code == 204:
650 # This response code is not allowed to have a non-empty body,
651 # and has an implicit length of zero instead of read-until-close.
652 # http://www.w3.org/Protocols/rfc2616/rfc2616-sec4.html#sec4.3
653 if is_chunked or content_length not in (None, 0):
654 raise httputil.HTTPInputError(
655 "Response with code %d should not have body" % code
656 )
657 content_length = 0
658
659 if is_chunked:
660 return self._read_chunked_body(delegate)
661 if content_length is not None:
662 return self._read_fixed_body(content_length, delegate)
663 if self.is_client:
664 return self._read_body_until_close(delegate)
665 return None
666
667 async def _read_fixed_body(
668 self, content_length: int, delegate: httputil.HTTPMessageDelegate
669 ) -> None:
670 while content_length > 0:
671 body = await self.stream.read_bytes(
672 min(self.params.chunk_size, content_length), partial=True
673 )
674 content_length -= len(body)
675 if not self._write_finished or self.is_client:
676 with _ExceptionLoggingContext(app_log):
677 ret = delegate.data_received(body)
678 if ret is not None:
679 await ret
680
681 async def _read_chunked_body(self, delegate: httputil.HTTPMessageDelegate) -> None:
682 # TODO: "chunk extensions" http://tools.ietf.org/html/rfc2616#section-3.6.1
683 total_size = 0
684 while True:
685 chunk_len_str = await self.stream.read_until(b"\r\n", max_bytes=64)
686 try:
687 chunk_len = parse_hex_int(native_str(chunk_len_str[:-2]))
688 except ValueError:
689 raise httputil.HTTPInputError("invalid chunk size")
690 if chunk_len == 0:
691 crlf = await self.stream.read_bytes(2)
692 if crlf != b"\r\n":
693 raise httputil.HTTPInputError(
694 "improperly terminated chunked request"
695 )
696 return
697 total_size += chunk_len
698 if total_size > self._max_body_size:
699 raise httputil.HTTPInputError("chunked body too large")
700 bytes_to_read = chunk_len
701 while bytes_to_read:
702 chunk = await self.stream.read_bytes(
703 min(bytes_to_read, self.params.chunk_size), partial=True
704 )
705 bytes_to_read -= len(chunk)
706 if not self._write_finished or self.is_client:
707 with _ExceptionLoggingContext(app_log):
708 ret = delegate.data_received(chunk)
709 if ret is not None:
710 await ret
711 # chunk ends with \r\n
712 crlf = await self.stream.read_bytes(2)
713 assert crlf == b"\r\n"
714
715 async def _read_body_until_close(
716 self, delegate: httputil.HTTPMessageDelegate
717 ) -> None:
718 # The body is terminated by the connection closing, so there is no
719 # length known in advance. Read incrementally so that max_body_size
720 # is enforced before an over-large body has been buffered, and so
721 # that the body is not limited by the stream's read buffer size.
722 total_size = 0
723 while True:
724 try:
725 body = await self.stream.read_bytes(
726 self.params.chunk_size, partial=True
727 )
728 except iostream.StreamClosedError:
729 # The connection closing is the normal end of this body.
730 return
731 total_size += len(body)
732 if total_size > self._max_body_size:
733 raise httputil.HTTPInputError("Body too long")
734 if not self._write_finished or self.is_client:
735 with _ExceptionLoggingContext(app_log):
736 ret = delegate.data_received(body)
737 if ret is not None:
738 await ret
739
740
741class _GzipMessageDelegate(httputil.HTTPMessageDelegate):
742 """Wraps an `HTTPMessageDelegate` to decode ``Content-Encoding: gzip``."""
743
744 def __init__(
745 self,
746 delegate: httputil.HTTPMessageDelegate,
747 chunk_size: int,
748 max_body_size: int,
749 ) -> None:
750 self._delegate = delegate
751 self._chunk_size = chunk_size
752 self._max_body_size = max_body_size
753 self._decompressed_body_size = 0
754 self._decompressor = None # type: Optional[GzipDecompressor]
755
756 def headers_received(
757 self,
758 start_line: Union[httputil.RequestStartLine, httputil.ResponseStartLine],
759 headers: httputil.HTTPHeaders,
760 ) -> Optional[Awaitable[None]]:
761 if headers.get("Content-Encoding", "").lower() == "gzip":
762 self._decompressor = GzipDecompressor()
763 # Downstream delegates will only see uncompressed data,
764 # so rename the content-encoding header.
765 # (but note that curl_httpclient doesn't do this).
766 headers.add("X-Consumed-Content-Encoding", headers["Content-Encoding"])
767 del headers["Content-Encoding"]
768 return self._delegate.headers_received(start_line, headers)
769
770 async def data_received(self, chunk: bytes) -> None:
771 if self._decompressor:
772 compressed_data = chunk
773 while compressed_data:
774 decompressed = self._decompressor.decompress(
775 compressed_data, self._chunk_size
776 )
777 if decompressed:
778 self._decompressed_body_size += len(decompressed)
779 if self._decompressed_body_size > self._max_body_size:
780 raise httputil.HTTPInputError("decompressed body too large")
781 ret = self._delegate.data_received(decompressed)
782 if ret is not None:
783 await ret
784 compressed_data = self._decompressor.unconsumed_tail
785 if compressed_data and not decompressed:
786 raise httputil.HTTPInputError(
787 "encountered unconsumed gzip data without making progress"
788 )
789 else:
790 ret = self._delegate.data_received(chunk)
791 if ret is not None:
792 await ret
793
794 def finish(self) -> None:
795 if self._decompressor is not None:
796 tail = self._decompressor.flush()
797 if tail:
798 # The tail should always be empty: decompress returned
799 # all that it can in data_received and the only
800 # purpose of the flush call is to detect errors such
801 # as truncated input. If we did legitimately get a new
802 # chunk at this point we'd need to change the
803 # interface to make finish() a coroutine.
804 raise ValueError(
805 "decompressor.flush returned data; possible truncated input"
806 )
807 return self._delegate.finish()
808
809 def on_connection_close(self) -> None:
810 return self._delegate.on_connection_close()
811
812
813class HTTP1ServerConnection:
814 """An HTTP/1.x server."""
815
816 def __init__(
817 self,
818 stream: iostream.IOStream,
819 params: Optional[HTTP1ConnectionParameters] = None,
820 context: Optional[object] = None,
821 ) -> None:
822 """
823 :arg stream: an `.IOStream`
824 :arg params: a `.HTTP1ConnectionParameters` or None
825 :arg context: an opaque application-defined object that is accessible
826 as ``connection.context``
827 """
828 self.stream = stream
829 if params is None:
830 params = HTTP1ConnectionParameters()
831 self.params = params
832 self.context = context
833 self._serving_future = None # type: Optional[Future[None]]
834
835 async def close(self) -> None:
836 """Closes the connection.
837
838 Returns a `.Future` that resolves after the serving loop has exited.
839 """
840 self.stream.close()
841 # Block until the serving loop is done, but ignore any exceptions
842 # (start_serving is already responsible for logging them).
843 assert self._serving_future is not None
844 try:
845 await self._serving_future
846 except Exception:
847 pass
848
849 def start_serving(self, delegate: httputil.HTTPServerConnectionDelegate) -> None:
850 """Starts serving requests on this connection.
851
852 :arg delegate: a `.HTTPServerConnectionDelegate`
853 """
854 assert isinstance(delegate, httputil.HTTPServerConnectionDelegate)
855 fut = gen.convert_yielded(self._server_request_loop(delegate))
856 self._serving_future = fut
857 # Register the future on the IOLoop so its errors get logged.
858 self.stream.io_loop.add_future(fut, lambda f: f.result())
859
860 async def _server_request_loop(
861 self, delegate: httputil.HTTPServerConnectionDelegate
862 ) -> None:
863 try:
864 while True:
865 conn = HTTP1Connection(self.stream, False, self.params, self.context)
866 request_delegate = delegate.start_request(self, conn)
867 try:
868 ret = await conn.read_response(request_delegate)
869 except (
870 iostream.StreamClosedError,
871 iostream.UnsatisfiableReadError,
872 asyncio.CancelledError,
873 ):
874 return
875 except _QuietException:
876 # This exception was already logged.
877 conn.close()
878 return
879 except Exception:
880 gen_log.error("Uncaught exception", exc_info=True)
881 conn.close()
882 return
883 if not ret:
884 return
885 await asyncio.sleep(0)
886 finally:
887 delegate.on_close(self)
888
889
890DIGITS = re.compile(r"[0-9]+")
891HEXDIGITS = re.compile(r"[0-9a-fA-F]+")
892
893
894def parse_int(s: str) -> int:
895 """Parse a non-negative integer from a string."""
896 if DIGITS.fullmatch(s) is None:
897 raise ValueError("not an integer: %r" % s)
898 return int(s)
899
900
901def parse_hex_int(s: str) -> int:
902 """Parse a non-negative hexadecimal integer from a string."""
903 if HEXDIGITS.fullmatch(s) is None:
904 raise ValueError("not a hexadecimal integer: %r" % s)
905 return int(s, 16)
906
907
908def is_transfer_encoding_chunked(headers: httputil.HTTPHeaders) -> bool:
909 """Returns true if the headers specify Transfer-Encoding: chunked.
910
911 Raise httputil.HTTPInputError if any other transfer encoding is used.
912 """
913 # Note that transfer-encoding is an area in which postel's law can lead
914 # us astray. If a proxy and a backend server are liberal in what they accept,
915 # but accept slightly different things, this can lead to mismatched framing
916 # and request smuggling issues. Therefore we are as strict as possible here
917 # (even technically going beyond the requirements of the RFCs: a value of
918 # ",chunked" is legal but doesn't appear in practice for legitimate traffic)
919 if "Transfer-Encoding" not in headers:
920 return False
921 if "Content-Length" in headers:
922 # Message cannot contain both Content-Length and
923 # Transfer-Encoding headers.
924 # http://tools.ietf.org/html/rfc7230#section-3.3.3
925 raise httputil.HTTPInputError(
926 "Message with both Transfer-Encoding and Content-Length"
927 )
928 if headers["Transfer-Encoding"].lower() == "chunked":
929 return True
930 # We do not support any transfer-encodings other than chunked, and we do not
931 # expect to add any support because the concept of transfer-encoding has
932 # been removed in HTTP/2.
933 raise httputil.HTTPInputError(
934 "Unsupported Transfer-Encoding %s" % headers["Transfer-Encoding"]
935 )