Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/tornado/http1connection.py: 13%

Shortcuts on this page

r m x   toggle line displays

j k   next/prev highlighted chunk

0   (zero) top of page

1   (one) first highlighted chunk

467 statements  

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 )