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

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

726 statements  

1# 

2# Copyright 2009 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"""Utility classes to write to and read from non-blocking files and sockets. 

17 

18Contents: 

19 

20* `BaseIOStream`: Generic interface for reading and writing. 

21* `IOStream`: Implementation of BaseIOStream using non-blocking sockets. 

22* `SSLIOStream`: SSL-aware version of IOStream. 

23* `PipeIOStream`: Pipe-based IOStream implementation. 

24""" 

25 

26import asyncio 

27import collections 

28import errno 

29import io 

30import numbers 

31import os 

32import socket 

33import ssl 

34import sys 

35import re 

36 

37from tornado.concurrent import Future, future_set_result_unless_cancelled 

38from tornado import ioloop 

39from tornado.log import gen_log 

40from tornado.netutil import ssl_wrap_socket, _client_ssl_defaults, _server_ssl_defaults 

41from tornado.util import errno_from_exception 

42 

43import typing 

44from typing import ( 

45 Union, 

46 Optional, 

47 Awaitable, 

48 Callable, 

49 Pattern, 

50 Any, 

51 Dict, 

52 TypeVar, 

53 Tuple, 

54) 

55from types import TracebackType 

56 

57if typing.TYPE_CHECKING: 

58 from typing import Deque, List, Type # noqa: F401 

59 

60_IOStreamType = TypeVar("_IOStreamType", bound="IOStream") 

61 

62# These errnos indicate that a connection has been abruptly terminated. 

63# They should be caught and handled less noisily than other errors. 

64_ERRNO_CONNRESET = (errno.ECONNRESET, errno.ECONNABORTED, errno.EPIPE, errno.ETIMEDOUT) 

65 

66if hasattr(errno, "WSAECONNRESET"): 

67 _ERRNO_CONNRESET += ( # type: ignore 

68 errno.WSAECONNRESET, # type: ignore 

69 errno.WSAECONNABORTED, # type: ignore 

70 errno.WSAETIMEDOUT, # type: ignore 

71 ) 

72 

73if sys.platform == "darwin": 

74 # OSX appears to have a race condition that causes send(2) to return 

75 # EPROTOTYPE if called while a socket is being torn down: 

76 # http://erickt.github.io/blog/2014/11/19/adventures-in-debugging-a-potential-osx-kernel-bug/ 

77 # Since the socket is being closed anyway, treat this as an ECONNRESET 

78 # instead of an unexpected error. 

79 _ERRNO_CONNRESET += (errno.EPROTOTYPE,) # type: ignore 

80 

81_WINDOWS = sys.platform.startswith("win") 

82 

83 

84class StreamClosedError(IOError): 

85 """Exception raised by `IOStream` methods when the stream is closed. 

86 

87 Note that the close callback is scheduled to run *after* other 

88 callbacks on the stream (to allow for buffered data to be processed), 

89 so you may see this error before you see the close callback. 

90 

91 The ``real_error`` attribute contains the underlying error that caused 

92 the stream to close (if any). 

93 

94 .. versionchanged:: 4.3 

95 Added the ``real_error`` attribute. 

96 """ 

97 

98 def __init__(self, real_error: Optional[BaseException] = None) -> None: 

99 super().__init__("Stream is closed") 

100 self.real_error = real_error 

101 

102 

103class UnsatisfiableReadError(Exception): 

104 """Exception raised when a read cannot be satisfied. 

105 

106 Raised by ``read_until`` and ``read_until_regex`` with a ``max_bytes`` 

107 argument. 

108 """ 

109 

110 pass 

111 

112 

113class StreamBufferFullError(Exception): 

114 """Exception raised by `IOStream` methods when the buffer is full.""" 

115 

116 

117class _StreamBuffer: 

118 """ 

119 A specialized buffer that tries to avoid copies when large pieces 

120 of data are encountered. 

121 """ 

122 

123 def __init__(self) -> None: 

124 # A sequence of (False, bytearray) and (True, memoryview) objects 

125 self._buffers = ( 

126 collections.deque() 

127 ) # type: Deque[Tuple[bool, Union[bytearray, memoryview]]] 

128 # Position in the first buffer 

129 self._first_pos = 0 

130 self._size = 0 

131 

132 def __len__(self) -> int: 

133 return self._size 

134 

135 # Data above this size will be appended separately instead 

136 # of extending an existing bytearray 

137 _large_buf_threshold = 2048 

138 

139 def append(self, data: Union[bytes, bytearray, memoryview]) -> None: 

140 """ 

141 Append the given piece of data (should be a buffer-compatible object). 

142 """ 

143 size = len(data) 

144 if size > self._large_buf_threshold: 

145 if not isinstance(data, memoryview): 

146 data = memoryview(data) 

147 self._buffers.append((True, data)) 

148 elif size > 0: 

149 if self._buffers: 

150 is_memview, b = self._buffers[-1] 

151 new_buf = is_memview or len(b) >= self._large_buf_threshold 

152 else: 

153 new_buf = True 

154 if new_buf: 

155 self._buffers.append((False, bytearray(data))) 

156 else: 

157 b += data # type: ignore 

158 

159 self._size += size 

160 

161 def peek(self, size: int) -> memoryview: 

162 """ 

163 Get a view over at most ``size`` bytes (possibly fewer) at the 

164 current buffer position. 

165 """ 

166 assert size > 0 

167 try: 

168 is_memview, b = self._buffers[0] 

169 except IndexError: 

170 return memoryview(b"") 

171 

172 pos = self._first_pos 

173 if is_memview: 

174 return typing.cast(memoryview, b[pos : pos + size]) 

175 else: 

176 return memoryview(b)[pos : pos + size] 

177 

178 def advance(self, size: int) -> None: 

179 """ 

180 Advance the current buffer position by ``size`` bytes. 

181 """ 

182 assert 0 < size <= self._size 

183 self._size -= size 

184 pos = self._first_pos 

185 

186 buffers = self._buffers 

187 while buffers and size > 0: 

188 is_large, b = buffers[0] 

189 b_remain = len(b) - size - pos 

190 if b_remain <= 0: 

191 buffers.popleft() 

192 size -= len(b) - pos 

193 pos = 0 

194 elif is_large: 

195 pos += size 

196 size = 0 

197 else: 

198 pos += size 

199 del typing.cast(bytearray, b)[:pos] 

200 pos = 0 

201 size = 0 

202 

203 assert size == 0 

204 self._first_pos = pos 

205 

206 

207class BaseIOStream: 

208 """A utility class to write to and read from a non-blocking file or socket. 

209 

210 We support a non-blocking ``write()`` and a family of ``read_*()`` 

211 methods. When the operation completes, the ``Awaitable`` will resolve 

212 with the data read (or ``None`` for ``write()``). All outstanding 

213 ``Awaitables`` will resolve with a `StreamClosedError` when the 

214 stream is closed; `.BaseIOStream.set_close_callback` can also be used 

215 to be notified of a closed stream. 

216 

217 When a stream is closed due to an error, the IOStream's ``error`` 

218 attribute contains the exception object. 

219 

220 Subclasses must implement `fileno`, `close_fd`, `write_to_fd`, 

221 `read_from_fd`, and optionally `get_fd_error`. 

222 

223 """ 

224 

225 def __init__( 

226 self, 

227 max_buffer_size: Optional[int] = None, 

228 read_chunk_size: Optional[int] = None, 

229 max_write_buffer_size: Optional[int] = None, 

230 ) -> None: 

231 """`BaseIOStream` constructor. 

232 

233 :arg max_buffer_size: Maximum amount of incoming data to buffer; 

234 defaults to 100MB. 

235 :arg read_chunk_size: Amount of data to read at one time from the 

236 underlying transport; defaults to 64KB. 

237 :arg max_write_buffer_size: Amount of outgoing data to buffer; 

238 defaults to unlimited. 

239 

240 .. versionchanged:: 4.0 

241 Add the ``max_write_buffer_size`` parameter. Changed default 

242 ``read_chunk_size`` to 64KB. 

243 .. versionchanged:: 5.0 

244 The ``io_loop`` argument (deprecated since version 4.1) has been 

245 removed. 

246 """ 

247 self.io_loop = ioloop.IOLoop.current() 

248 self.max_buffer_size = max_buffer_size or 104857600 

249 # A chunk size that is too close to max_buffer_size can cause 

250 # spurious failures. 

251 self.read_chunk_size = min(read_chunk_size or 65536, self.max_buffer_size // 2) 

252 self.max_write_buffer_size = max_write_buffer_size 

253 self.error = None # type: Optional[BaseException] 

254 self._read_buffer = bytearray() 

255 self._read_buffer_size = 0 

256 self._user_read_buffer = False 

257 self._after_user_read_buffer = None # type: Optional[bytearray] 

258 self._write_buffer = _StreamBuffer() 

259 self._total_write_index = 0 

260 self._total_write_done_index = 0 

261 self._read_delimiter = None # type: Optional[bytes] 

262 self._read_regex = None # type: Optional[Pattern] 

263 self._read_max_bytes = None # type: Optional[int] 

264 self._read_bytes = None # type: Optional[int] 

265 self._read_partial = False 

266 self._read_until_close = False 

267 self._read_future = None # type: Optional[Future] 

268 self._write_futures = ( 

269 collections.deque() 

270 ) # type: Deque[Tuple[int, Future[None]]] 

271 self._close_callback = None # type: Optional[Callable[[], None]] 

272 self._connect_future = None # type: Optional[Future[IOStream]] 

273 # _ssl_connect_future should be defined in SSLIOStream 

274 # but it's here so we can clean it up in _signal_closed 

275 # TODO: refactor that so subclasses can add additional futures 

276 # to be cancelled. 

277 self._ssl_connect_future = None # type: Optional[Future[SSLIOStream]] 

278 self._connecting = False 

279 self._state = None # type: Optional[int] 

280 self._closed = False 

281 

282 def fileno(self) -> Union[int, ioloop._Selectable]: 

283 """Returns the file descriptor for this stream.""" 

284 raise NotImplementedError() 

285 

286 def close_fd(self) -> None: 

287 """Closes the file underlying this stream. 

288 

289 ``close_fd`` is called by `BaseIOStream` and should not be called 

290 elsewhere; other users should call `close` instead. 

291 """ 

292 raise NotImplementedError() 

293 

294 def write_to_fd(self, data: memoryview) -> int: 

295 """Attempts to write ``data`` to the underlying file. 

296 

297 Returns the number of bytes written. 

298 """ 

299 raise NotImplementedError() 

300 

301 def read_from_fd(self, buf: Union[bytearray, memoryview]) -> Optional[int]: 

302 """Attempts to read from the underlying file. 

303 

304 Reads up to ``len(buf)`` bytes, storing them in the buffer. 

305 Returns the number of bytes read. Returns None if there was 

306 nothing to read (the socket returned `~errno.EWOULDBLOCK` or 

307 equivalent), and zero on EOF. 

308 

309 .. versionchanged:: 5.0 

310 

311 Interface redesigned to take a buffer and return a number 

312 of bytes instead of a freshly-allocated object. 

313 """ 

314 raise NotImplementedError() 

315 

316 def get_fd_error(self) -> Optional[Exception]: 

317 """Returns information about any error on the underlying file. 

318 

319 This method is called after the `.IOLoop` has signaled an error on the 

320 file descriptor, and should return an Exception (such as `socket.error` 

321 with additional information, or None if no such information is 

322 available. 

323 """ 

324 return None 

325 

326 def read_until_regex( 

327 self, regex: bytes, max_bytes: Optional[int] = None 

328 ) -> Awaitable[bytes]: 

329 """Asynchronously read until we have matched the given regex. 

330 

331 The result includes the data that matches the regex and anything 

332 that came before it. 

333 

334 If ``max_bytes`` is not None, the connection will be closed 

335 if more than ``max_bytes`` bytes have been read and the regex is 

336 not satisfied. 

337 

338 .. versionchanged:: 4.0 

339 Added the ``max_bytes`` argument. The ``callback`` argument is 

340 now optional and a `.Future` will be returned if it is omitted. 

341 

342 .. versionchanged:: 6.0 

343 

344 The ``callback`` argument was removed. Use the returned 

345 `.Future` instead. 

346 

347 """ 

348 future = self._start_read() 

349 self._read_regex = re.compile(regex) 

350 self._read_max_bytes = max_bytes 

351 try: 

352 self._try_inline_read() 

353 except UnsatisfiableReadError as e: 

354 # Handle this the same way as in _handle_events. 

355 gen_log.info("Unsatisfiable read, closing connection: %s" % e) 

356 self.close(exc_info=e) 

357 return future 

358 except: 

359 # Ensure that the future doesn't log an error because its 

360 # failure was never examined. 

361 future.add_done_callback(lambda f: f.exception()) 

362 raise 

363 return future 

364 

365 def read_until( 

366 self, delimiter: bytes, max_bytes: Optional[int] = None 

367 ) -> Awaitable[bytes]: 

368 """Asynchronously read until we have found the given delimiter. 

369 

370 The result includes all the data read including the delimiter. 

371 

372 If ``max_bytes`` is not None, the connection will be closed 

373 if more than ``max_bytes`` bytes have been read and the delimiter 

374 is not found. 

375 

376 .. versionchanged:: 4.0 

377 Added the ``max_bytes`` argument. The ``callback`` argument is 

378 now optional and a `.Future` will be returned if it is omitted. 

379 

380 .. versionchanged:: 6.0 

381 

382 The ``callback`` argument was removed. Use the returned 

383 `.Future` instead. 

384 """ 

385 future = self._start_read() 

386 self._read_delimiter = delimiter 

387 self._read_max_bytes = max_bytes 

388 try: 

389 self._try_inline_read() 

390 except UnsatisfiableReadError as e: 

391 # Handle this the same way as in _handle_events. 

392 gen_log.info("Unsatisfiable read, closing connection: %s" % e) 

393 self.close(exc_info=e) 

394 return future 

395 except: 

396 future.add_done_callback(lambda f: f.exception()) 

397 raise 

398 return future 

399 

400 def read_bytes(self, num_bytes: int, partial: bool = False) -> Awaitable[bytes]: 

401 """Asynchronously read a number of bytes. 

402 

403 If ``partial`` is true, data is returned as soon as we have 

404 any bytes to return (but never more than ``num_bytes``) 

405 

406 .. versionchanged:: 4.0 

407 Added the ``partial`` argument. The callback argument is now 

408 optional and a `.Future` will be returned if it is omitted. 

409 

410 .. versionchanged:: 6.0 

411 

412 The ``callback`` and ``streaming_callback`` arguments have 

413 been removed. Use the returned `.Future` (and 

414 ``partial=True`` for ``streaming_callback``) instead. 

415 

416 """ 

417 future = self._start_read() 

418 assert isinstance(num_bytes, numbers.Integral) 

419 self._read_bytes = num_bytes 

420 self._read_partial = partial 

421 try: 

422 self._try_inline_read() 

423 except: 

424 future.add_done_callback(lambda f: f.exception()) 

425 raise 

426 return future 

427 

428 def read_into(self, buf: bytearray, partial: bool = False) -> Awaitable[int]: 

429 """Asynchronously read a number of bytes. 

430 

431 ``buf`` must be a writable buffer into which data will be read. 

432 

433 If ``partial`` is true, the callback is run as soon as any bytes 

434 have been read. Otherwise, it is run when the ``buf`` has been 

435 entirely filled with read data. 

436 

437 .. versionadded:: 5.0 

438 

439 .. versionchanged:: 6.0 

440 

441 The ``callback`` argument was removed. Use the returned 

442 `.Future` instead. 

443 

444 """ 

445 future = self._start_read() 

446 

447 # First copy data already in read buffer 

448 available_bytes = self._read_buffer_size 

449 n = len(buf) 

450 if available_bytes >= n: 

451 buf[:] = memoryview(self._read_buffer)[:n] 

452 del self._read_buffer[:n] 

453 self._after_user_read_buffer = self._read_buffer 

454 elif available_bytes > 0: 

455 buf[:available_bytes] = memoryview(self._read_buffer)[:] 

456 

457 # Set up the supplied buffer as our temporary read buffer. 

458 # The original (if it had any data remaining) has been 

459 # saved for later. 

460 self._user_read_buffer = True 

461 self._read_buffer = buf 

462 self._read_buffer_size = available_bytes 

463 self._read_bytes = n 

464 self._read_partial = partial 

465 

466 try: 

467 self._try_inline_read() 

468 except: 

469 future.add_done_callback(lambda f: f.exception()) 

470 raise 

471 return future 

472 

473 def read_until_close(self) -> Awaitable[bytes]: 

474 """Asynchronously reads all data from the socket until it is closed. 

475 

476 This will buffer all available data until ``max_buffer_size`` 

477 is reached. If flow control or cancellation are desired, use a 

478 loop with `read_bytes(partial=True) <.read_bytes>` instead. 

479 

480 .. versionchanged:: 4.0 

481 The callback argument is now optional and a `.Future` will 

482 be returned if it is omitted. 

483 

484 .. versionchanged:: 6.0 

485 

486 The ``callback`` and ``streaming_callback`` arguments have 

487 been removed. Use the returned `.Future` (and `read_bytes` 

488 with ``partial=True`` for ``streaming_callback``) instead. 

489 

490 """ 

491 future = self._start_read() 

492 if self.closed(): 

493 self._finish_read(self._read_buffer_size) 

494 return future 

495 self._read_until_close = True 

496 try: 

497 self._try_inline_read() 

498 except: 

499 future.add_done_callback(lambda f: f.exception()) 

500 raise 

501 return future 

502 

503 def write(self, data: Union[bytes, memoryview]) -> "Future[None]": 

504 """Asynchronously write the given data to this stream. 

505 

506 This method returns a `.Future` that resolves (with a result 

507 of ``None``) when the write has been completed. 

508 

509 The ``data`` argument may be of type `bytes` or `memoryview`. 

510 

511 .. versionchanged:: 4.0 

512 Now returns a `.Future` if no callback is given. 

513 

514 .. versionchanged:: 4.5 

515 Added support for `memoryview` arguments. 

516 

517 .. versionchanged:: 6.0 

518 

519 The ``callback`` argument was removed. Use the returned 

520 `.Future` instead. 

521 

522 """ 

523 self._check_closed() 

524 if data: 

525 if isinstance(data, memoryview): 

526 # Make sure that ``len(data) == data.nbytes`` 

527 data = memoryview(data).cast("B") 

528 if ( 

529 self.max_write_buffer_size is not None 

530 and len(self._write_buffer) + len(data) > self.max_write_buffer_size 

531 ): 

532 raise StreamBufferFullError("Reached maximum write buffer size") 

533 self._write_buffer.append(data) 

534 self._total_write_index += len(data) 

535 future = Future() # type: Future[None] 

536 future.add_done_callback(lambda f: f.exception()) 

537 self._write_futures.append((self._total_write_index, future)) 

538 if not self._connecting: 

539 self._handle_write() 

540 if self._write_buffer: 

541 self._add_io_state(self.io_loop.WRITE) 

542 self._maybe_add_error_listener() 

543 return future 

544 

545 def set_close_callback(self, callback: Optional[Callable[[], None]]) -> None: 

546 """Call the given callback when the stream is closed. 

547 

548 This mostly is not necessary for applications that use the 

549 `.Future` interface; all outstanding ``Futures`` will resolve 

550 with a `StreamClosedError` when the stream is closed. However, 

551 it is still useful as a way to signal that the stream has been 

552 closed while no other read or write is in progress. 

553 

554 Unlike other callback-based interfaces, ``set_close_callback`` 

555 was not removed in Tornado 6.0. 

556 """ 

557 self._close_callback = callback 

558 self._maybe_add_error_listener() 

559 

560 def close( 

561 self, 

562 exc_info: Union[ 

563 None, 

564 bool, 

565 BaseException, 

566 Tuple[ 

567 "Optional[Type[BaseException]]", 

568 Optional[BaseException], 

569 Optional[TracebackType], 

570 ], 

571 ] = False, 

572 ) -> None: 

573 """Close this stream. 

574 

575 If ``exc_info`` is true, set the ``error`` attribute to the current 

576 exception from `sys.exc_info` (or if ``exc_info`` is a tuple, 

577 use that instead of `sys.exc_info`). 

578 """ 

579 if not self.closed(): 

580 if exc_info: 

581 if isinstance(exc_info, tuple): 

582 self.error = exc_info[1] 

583 elif isinstance(exc_info, BaseException): 

584 self.error = exc_info 

585 else: 

586 exc_info = sys.exc_info() 

587 if any(exc_info): 

588 self.error = exc_info[1] 

589 if self._read_until_close: 

590 self._read_until_close = False 

591 if self.error is None or self._is_connreset(self.error): 

592 # A connection reset is treated as a normal close 

593 # throughout this class (on some platforms, notably 

594 # windows, a peer that closes cleanly may still be 

595 # reported as a reset), so deliver the buffered data as 

596 # the result of the read. 

597 self._finish_read(self._read_buffer_size) 

598 # Otherwise the stream is closing because of a real error, so 

599 # leave the read future pending for _signal_closed() to fail 

600 # with StreamClosedError(real_error=self.error). Resolving it 

601 # with the buffered data would report a truncated result as if 

602 # it were complete. 

603 elif self._read_future is not None: 

604 # resolve reads that are pending and ready to complete 

605 try: 

606 pos = self._find_read_pos() 

607 except UnsatisfiableReadError: 

608 pass 

609 else: 

610 if pos is not None: 

611 self._read_from_buffer(pos) 

612 if self._state is not None: 

613 self.io_loop.remove_handler(self.fileno()) 

614 self._state = None 

615 self.close_fd() 

616 self._closed = True 

617 self._signal_closed() 

618 

619 def _signal_closed(self) -> None: 

620 futures = [] # type: List[Future] 

621 if self._read_future is not None: 

622 futures.append(self._read_future) 

623 self._read_future = None 

624 futures += [future for _, future in self._write_futures] 

625 self._write_futures.clear() 

626 if self._connect_future is not None: 

627 futures.append(self._connect_future) 

628 self._connect_future = None 

629 for future in futures: 

630 if not future.done(): 

631 future.set_exception(StreamClosedError(real_error=self.error)) 

632 # Reference the exception to silence warnings. Annoyingly, 

633 # this raises if the future was cancelled, but just 

634 # returns any other error. 

635 try: 

636 future.exception() 

637 except asyncio.CancelledError: 

638 pass 

639 if self._ssl_connect_future is not None: 

640 # _ssl_connect_future expects to see the real exception (typically 

641 # an ssl.SSLError), not just StreamClosedError. 

642 if not self._ssl_connect_future.done(): 

643 if self.error is not None: 

644 self._ssl_connect_future.set_exception(self.error) 

645 else: 

646 self._ssl_connect_future.set_exception(StreamClosedError()) 

647 self._ssl_connect_future.exception() 

648 self._ssl_connect_future = None 

649 if self._close_callback is not None: 

650 cb = self._close_callback 

651 self._close_callback = None 

652 self.io_loop.add_callback(cb) 

653 # Clear the buffers so they can be cleared immediately even 

654 # if the IOStream object is kept alive by a reference cycle. 

655 # TODO: Clear the read buffer too; it currently breaks some tests. 

656 self._write_buffer = None # type: ignore 

657 

658 def reading(self) -> bool: 

659 """Returns ``True`` if we are currently reading from the stream.""" 

660 return self._read_future is not None 

661 

662 def writing(self) -> bool: 

663 """Returns ``True`` if we are currently writing to the stream.""" 

664 return bool(self._write_buffer) 

665 

666 def closed(self) -> bool: 

667 """Returns ``True`` if the stream has been closed.""" 

668 return self._closed 

669 

670 def set_nodelay(self, value: bool) -> None: 

671 """Sets the no-delay flag for this stream. 

672 

673 By default, data written to TCP streams may be held for a time 

674 to make the most efficient use of bandwidth (according to 

675 Nagle's algorithm). The no-delay flag requests that data be 

676 written as soon as possible, even if doing so would consume 

677 additional bandwidth. 

678 

679 This flag is currently defined only for TCP-based ``IOStreams``. 

680 

681 .. versionadded:: 3.1 

682 """ 

683 pass 

684 

685 def _handle_connect(self) -> None: 

686 raise NotImplementedError() 

687 

688 def _handle_events(self, fd: Union[int, ioloop._Selectable], events: int) -> None: 

689 if self.closed(): 

690 gen_log.warning("Got events for closed stream %s", fd) 

691 return 

692 try: 

693 if self._connecting: 

694 # Most IOLoops will report a write failed connect 

695 # with the WRITE event, but SelectIOLoop reports a 

696 # READ as well so we must check for connecting before 

697 # either. 

698 self._handle_connect() 

699 if self.closed(): 

700 return 

701 if events & self.io_loop.READ: 

702 self._handle_read() 

703 if self.closed(): 

704 return 

705 if events & self.io_loop.WRITE: 

706 self._handle_write() 

707 if self.closed(): 

708 return 

709 if events & self.io_loop.ERROR: 

710 self.error = self.get_fd_error() 

711 # We may have queued up a user callback in _handle_read or 

712 # _handle_write, so don't close the IOStream until those 

713 # callbacks have had a chance to run. 

714 self.io_loop.add_callback(self.close) 

715 return 

716 state = self.io_loop.ERROR 

717 if self.reading(): 

718 state |= self.io_loop.READ 

719 if self.writing(): 

720 state |= self.io_loop.WRITE 

721 if state == self.io_loop.ERROR and self._read_buffer_size == 0: 

722 # If the connection is idle, listen for reads too so 

723 # we can tell if the connection is closed. If there is 

724 # data in the read buffer we won't run the close callback 

725 # yet anyway, so we don't need to listen in this case. 

726 state |= self.io_loop.READ 

727 if state != self._state: 

728 assert ( 

729 self._state is not None 

730 ), "shouldn't happen: _handle_events without self._state" 

731 self._state = state 

732 self.io_loop.update_handler(self.fileno(), self._state) 

733 except UnsatisfiableReadError as e: 

734 gen_log.info("Unsatisfiable read, closing connection: %s" % e) 

735 self.close(exc_info=e) 

736 except Exception as e: 

737 gen_log.error("Uncaught exception, closing connection.", exc_info=True) 

738 self.close(exc_info=e) 

739 raise 

740 

741 def _read_to_buffer_loop(self) -> Optional[int]: 

742 # This method is called from _handle_read and _try_inline_read. 

743 if self._read_bytes is not None: 

744 target_bytes = self._read_bytes # type: Optional[int] 

745 elif self._read_max_bytes is not None: 

746 target_bytes = self._read_max_bytes 

747 elif self.reading(): 

748 # For read_until without max_bytes, or 

749 # read_until_close, read as much as we can before 

750 # scanning for the delimiter. 

751 target_bytes = None 

752 else: 

753 target_bytes = 0 

754 next_find_pos = 0 

755 while not self.closed(): 

756 # Read from the socket until we get EWOULDBLOCK or equivalent. 

757 # SSL sockets do some internal buffering, and if the data is 

758 # sitting in the SSL object's buffer select() and friends 

759 # can't see it; the only way to find out if it's there is to 

760 # try to read it. 

761 if self._read_to_buffer() == 0: 

762 break 

763 

764 # If we've read all the bytes we can use, break out of 

765 # this loop. 

766 

767 # If we've reached target_bytes, we know we're done. 

768 if target_bytes is not None and self._read_buffer_size >= target_bytes: 

769 break 

770 

771 # Otherwise, we need to call the more expensive find_read_pos. 

772 # It's inefficient to do this on every read, so instead 

773 # do it on the first read and whenever the read buffer 

774 # size has doubled. 

775 if self._read_buffer_size >= next_find_pos: 

776 pos = self._find_read_pos() 

777 if pos is not None: 

778 return pos 

779 next_find_pos = self._read_buffer_size * 2 

780 return self._find_read_pos() 

781 

782 def _handle_read(self) -> None: 

783 try: 

784 pos = self._read_to_buffer_loop() 

785 except UnsatisfiableReadError: 

786 raise 

787 except asyncio.CancelledError: 

788 raise 

789 except Exception as e: 

790 gen_log.warning("error on read: %s" % e) 

791 self.close(exc_info=e) 

792 return 

793 if pos is not None: 

794 self._read_from_buffer(pos) 

795 

796 def _start_read(self) -> Future: 

797 if self._read_future is not None: 

798 # It is an error to start a read while a prior read is unresolved. 

799 # However, if the prior read is unresolved because the stream was 

800 # closed without satisfying it, it's better to raise 

801 # StreamClosedError instead of AssertionError. In particular, this 

802 # situation occurs in harmless situations in http1connection.py and 

803 # an AssertionError would be logged noisily. 

804 # 

805 # On the other hand, it is legal to start a new read while the 

806 # stream is closed, in case the read can be satisfied from the 

807 # read buffer. So we only want to check the closed status of the 

808 # stream if we need to decide what kind of error to raise for 

809 # "already reading". 

810 # 

811 # These conditions have proven difficult to test; we have no 

812 # unittests that reliably verify this behavior so be careful 

813 # when making changes here. See #2651 and #2719. 

814 self._check_closed() 

815 assert self._read_future is None, "Already reading" 

816 self._read_future = Future() 

817 return self._read_future 

818 

819 def _finish_read(self, size: int) -> None: 

820 if self._user_read_buffer: 

821 self._read_buffer = self._after_user_read_buffer or bytearray() 

822 self._after_user_read_buffer = None 

823 self._read_buffer_size = len(self._read_buffer) 

824 self._user_read_buffer = False 

825 result = size # type: Union[int, bytes] 

826 else: 

827 result = self._consume(size) 

828 if self._read_future is not None: 

829 future = self._read_future 

830 self._read_future = None 

831 future_set_result_unless_cancelled(future, result) 

832 self._maybe_add_error_listener() 

833 

834 def _try_inline_read(self) -> None: 

835 """Attempt to complete the current read operation from buffered data. 

836 

837 If the read can be completed without blocking, schedules the 

838 read callback on the next IOLoop iteration; otherwise starts 

839 listening for reads on the socket. 

840 """ 

841 # See if we've already got the data from a previous read 

842 pos = self._find_read_pos() 

843 if pos is not None: 

844 self._read_from_buffer(pos) 

845 return 

846 self._check_closed() 

847 pos = self._read_to_buffer_loop() 

848 if pos is not None: 

849 self._read_from_buffer(pos) 

850 return 

851 # We couldn't satisfy the read inline, so make sure we're 

852 # listening for new data unless the stream is closed. 

853 if not self.closed(): 

854 self._add_io_state(ioloop.IOLoop.READ) 

855 

856 def _read_to_buffer(self) -> Optional[int]: 

857 """Reads from the socket and appends the result to the read buffer. 

858 

859 Returns the number of bytes read. Returns 0 if there is nothing 

860 to read (i.e. the read returns EWOULDBLOCK or equivalent). On 

861 error closes the socket and raises an exception. 

862 """ 

863 try: 

864 while True: 

865 try: 

866 if self._user_read_buffer: 

867 buf = memoryview(self._read_buffer)[ 

868 self._read_buffer_size : 

869 ] # type: Union[memoryview, bytearray] 

870 else: 

871 buf = bytearray(self.read_chunk_size) 

872 bytes_read = self.read_from_fd(buf) 

873 except OSError as e: 

874 # ssl.SSLError is a subclass of socket.error 

875 if self._is_connreset(e): 

876 # Treat ECONNRESET as a connection close rather than 

877 # an error to minimize log spam (the exception will 

878 # be available on self.error for apps that care). 

879 self.close(exc_info=e) 

880 return None 

881 self.close(exc_info=e) 

882 raise 

883 break 

884 if bytes_read is None: 

885 return 0 

886 elif bytes_read == 0: 

887 self.close() 

888 return 0 

889 if not self._user_read_buffer: 

890 self._read_buffer += memoryview(buf)[:bytes_read] 

891 self._read_buffer_size += bytes_read 

892 finally: 

893 # Break the reference to buf so we don't waste a chunk's worth of 

894 # memory in case an exception hangs on to our stack frame. 

895 del buf 

896 if self._read_buffer_size > self.max_buffer_size: 

897 gen_log.error("Reached maximum read buffer size") 

898 buffer_full_error = StreamBufferFullError( 

899 "Reached maximum read buffer size" 

900 ) 

901 self.close(exc_info=buffer_full_error) 

902 raise buffer_full_error 

903 return bytes_read 

904 

905 def _read_from_buffer(self, pos: int) -> None: 

906 """Attempts to complete the currently-pending read from the buffer. 

907 

908 The argument is either a position in the read buffer or None, 

909 as returned by _find_read_pos. 

910 """ 

911 self._read_bytes = self._read_delimiter = self._read_regex = None 

912 self._read_partial = False 

913 self._finish_read(pos) 

914 

915 def _find_read_pos(self) -> Optional[int]: 

916 """Attempts to find a position in the read buffer that satisfies 

917 the currently-pending read. 

918 

919 Returns a position in the buffer if the current read can be satisfied, 

920 or None if it cannot. 

921 """ 

922 if self._read_bytes is not None and ( 

923 self._read_buffer_size >= self._read_bytes 

924 or (self._read_partial and self._read_buffer_size > 0) 

925 ): 

926 num_bytes = min(self._read_bytes, self._read_buffer_size) 

927 return num_bytes 

928 elif self._read_delimiter is not None: 

929 # Multi-byte delimiters (e.g. '\r\n') may straddle two 

930 # chunks in the read buffer, so we can't easily find them 

931 # without collapsing the buffer. However, since protocols 

932 # using delimited reads (as opposed to reads of a known 

933 # length) tend to be "line" oriented, the delimiter is likely 

934 # to be in the first few chunks. Merge the buffer gradually 

935 # since large merges are relatively expensive and get undone in 

936 # _consume(). 

937 if self._read_buffer: 

938 loc = self._read_buffer.find(self._read_delimiter) 

939 if loc != -1: 

940 delimiter_len = len(self._read_delimiter) 

941 self._check_max_bytes(self._read_delimiter, loc + delimiter_len) 

942 return loc + delimiter_len 

943 self._check_max_bytes(self._read_delimiter, self._read_buffer_size) 

944 elif self._read_regex is not None: 

945 if self._read_buffer: 

946 m = self._read_regex.search(self._read_buffer) 

947 if m is not None: 

948 loc = m.end() 

949 self._check_max_bytes(self._read_regex, loc) 

950 return loc 

951 self._check_max_bytes(self._read_regex, self._read_buffer_size) 

952 return None 

953 

954 def _check_max_bytes(self, delimiter: Union[bytes, Pattern], size: int) -> None: 

955 if self._read_max_bytes is not None and size > self._read_max_bytes: 

956 raise UnsatisfiableReadError( 

957 "delimiter %r not found within %d bytes" 

958 % (delimiter, self._read_max_bytes) 

959 ) 

960 

961 def _handle_write(self) -> None: 

962 while True: 

963 size = len(self._write_buffer) 

964 if not size: 

965 break 

966 assert size > 0 

967 try: 

968 if _WINDOWS: 

969 # On windows, socket.send blows up if given a 

970 # write buffer that's too large, instead of just 

971 # returning the number of bytes it was able to 

972 # process. Therefore we must not call socket.send 

973 # with more than 128KB at a time. 

974 size = 128 * 1024 

975 

976 num_bytes = self.write_to_fd(self._write_buffer.peek(size)) 

977 if num_bytes == 0: 

978 break 

979 self._write_buffer.advance(num_bytes) 

980 self._total_write_done_index += num_bytes 

981 except BlockingIOError: 

982 break 

983 except OSError as e: 

984 if not self._is_connreset(e): 

985 # Broken pipe errors are usually caused by connection 

986 # reset, and its better to not log EPIPE errors to 

987 # minimize log spam 

988 gen_log.warning("Write error on %s: %s", self.fileno(), e) 

989 self.close(exc_info=e) 

990 return 

991 

992 while self._write_futures: 

993 index, future = self._write_futures[0] 

994 if index > self._total_write_done_index: 

995 break 

996 self._write_futures.popleft() 

997 future_set_result_unless_cancelled(future, None) 

998 

999 def _consume(self, loc: int) -> bytes: 

1000 # Consume loc bytes from the read buffer and return them 

1001 if loc == 0: 

1002 return b"" 

1003 assert loc <= self._read_buffer_size 

1004 # Slice the bytearray buffer into bytes, without intermediate copying 

1005 b = (memoryview(self._read_buffer)[:loc]).tobytes() 

1006 self._read_buffer_size -= loc 

1007 del self._read_buffer[:loc] 

1008 return b 

1009 

1010 def _check_closed(self) -> None: 

1011 if self.closed(): 

1012 raise StreamClosedError(real_error=self.error) 

1013 

1014 def _maybe_add_error_listener(self) -> None: 

1015 # This method is part of an optimization: to detect a connection that 

1016 # is closed when we're not actively reading or writing, we must listen 

1017 # for read events. However, it is inefficient to do this when the 

1018 # connection is first established because we are going to read or write 

1019 # immediately anyway. Instead, we insert checks at various times to 

1020 # see if the connection is idle and add the read listener then. 

1021 if self._state is None or self._state == ioloop.IOLoop.ERROR: 

1022 if ( 

1023 not self.closed() 

1024 and self._read_buffer_size == 0 

1025 and self._close_callback is not None 

1026 ): 

1027 self._add_io_state(ioloop.IOLoop.READ) 

1028 

1029 def _add_io_state(self, state: int) -> None: 

1030 """Adds `state` (IOLoop.{READ,WRITE} flags) to our event handler. 

1031 

1032 Implementation notes: Reads and writes have a fast path and a 

1033 slow path. The fast path reads synchronously from socket 

1034 buffers, while the slow path uses `_add_io_state` to schedule 

1035 an IOLoop callback. 

1036 

1037 To detect closed connections, we must have called 

1038 `_add_io_state` at some point, but we want to delay this as 

1039 much as possible so we don't have to set an `IOLoop.ERROR` 

1040 listener that will be overwritten by the next slow-path 

1041 operation. If a sequence of fast-path ops do not end in a 

1042 slow-path op, (e.g. for an @asynchronous long-poll request), 

1043 we must add the error handler. 

1044 

1045 TODO: reevaluate this now that callbacks are gone. 

1046 

1047 """ 

1048 if self.closed(): 

1049 # connection has been closed, so there can be no future events 

1050 return 

1051 if self._state is None: 

1052 self._state = ioloop.IOLoop.ERROR | state 

1053 self.io_loop.add_handler(self.fileno(), self._handle_events, self._state) 

1054 elif not self._state & state: 

1055 self._state = self._state | state 

1056 self.io_loop.update_handler(self.fileno(), self._state) 

1057 

1058 def _is_connreset(self, exc: BaseException) -> bool: 

1059 """Return ``True`` if exc is ECONNRESET or equivalent. 

1060 

1061 May be overridden in subclasses. 

1062 """ 

1063 return ( 

1064 isinstance(exc, (socket.error, IOError)) 

1065 and errno_from_exception(exc) in _ERRNO_CONNRESET 

1066 ) 

1067 

1068 

1069class IOStream(BaseIOStream): 

1070 r"""Socket-based `IOStream` implementation. 

1071 

1072 This class supports the read and write methods from `BaseIOStream` 

1073 plus a `connect` method. 

1074 

1075 The ``socket`` parameter may either be connected or unconnected. 

1076 For server operations the socket is the result of calling 

1077 `socket.accept <socket.socket.accept>`. For client operations the 

1078 socket is created with `socket.socket`, and may either be 

1079 connected before passing it to the `IOStream` or connected with 

1080 `IOStream.connect`. 

1081 

1082 A very simple (and broken) HTTP client using this class: 

1083 

1084 .. testcode:: 

1085 

1086 import socket 

1087 import tornado 

1088 

1089 async def main(): 

1090 s = socket.socket(socket.AF_INET, socket.SOCK_STREAM, 0) 

1091 stream = tornado.iostream.IOStream(s) 

1092 await stream.connect(("friendfeed.com", 80)) 

1093 await stream.write(b"GET / HTTP/1.0\r\nHost: friendfeed.com\r\n\r\n") 

1094 header_data = await stream.read_until(b"\r\n\r\n") 

1095 headers = {} 

1096 for line in header_data.split(b"\r\n"): 

1097 parts = line.split(b":") 

1098 if len(parts) == 2: 

1099 headers[parts[0].strip()] = parts[1].strip() 

1100 body_data = await stream.read_bytes(int(headers[b"Content-Length"])) 

1101 print(body_data) 

1102 stream.close() 

1103 

1104 if __name__ == '__main__': 

1105 asyncio.run(main()) 

1106 

1107 """ 

1108 

1109 def __init__(self, socket: socket.socket, *args: Any, **kwargs: Any) -> None: 

1110 self.socket = socket 

1111 self.socket.setblocking(False) 

1112 super().__init__(*args, **kwargs) 

1113 

1114 def fileno(self) -> Union[int, ioloop._Selectable]: 

1115 return self.socket 

1116 

1117 def close_fd(self) -> None: 

1118 self.socket.close() 

1119 self.socket = None # type: ignore 

1120 

1121 def get_fd_error(self) -> Optional[Exception]: 

1122 errno = self.socket.getsockopt(socket.SOL_SOCKET, socket.SO_ERROR) 

1123 return socket.error(errno, os.strerror(errno)) 

1124 

1125 def read_from_fd(self, buf: Union[bytearray, memoryview]) -> Optional[int]: 

1126 try: 

1127 return self.socket.recv_into(buf, len(buf)) 

1128 except BlockingIOError: 

1129 return None 

1130 finally: 

1131 del buf 

1132 

1133 def write_to_fd(self, data: memoryview) -> int: 

1134 try: 

1135 return self.socket.send(data) # type: ignore 

1136 finally: 

1137 # Avoid keeping to data, which can be a memoryview. 

1138 # See https://github.com/tornadoweb/tornado/pull/2008 

1139 del data 

1140 

1141 def connect( 

1142 self: _IOStreamType, address: Any, server_hostname: Optional[str] = None 

1143 ) -> "Future[_IOStreamType]": 

1144 """Connects the socket to a remote address without blocking. 

1145 

1146 May only be called if the socket passed to the constructor was 

1147 not previously connected. The address parameter is in the 

1148 same format as for `socket.connect <socket.socket.connect>` for 

1149 the type of socket passed to the IOStream constructor, 

1150 e.g. an ``(ip, port)`` tuple. Hostnames are accepted here, 

1151 but will be resolved synchronously and block the IOLoop. 

1152 If you have a hostname instead of an IP address, the `.TCPClient` 

1153 class is recommended instead of calling this method directly. 

1154 `.TCPClient` will do asynchronous DNS resolution and handle 

1155 both IPv4 and IPv6. 

1156 

1157 If ``callback`` is specified, it will be called with no 

1158 arguments when the connection is completed; if not this method 

1159 returns a `.Future` (whose result after a successful 

1160 connection will be the stream itself). 

1161 

1162 In SSL mode, the ``server_hostname`` parameter will be used 

1163 for certificate validation (unless disabled in the 

1164 ``ssl_options``) and SNI. 

1165 

1166 Note that it is safe to call `IOStream.write 

1167 <BaseIOStream.write>` while the connection is pending, in 

1168 which case the data will be written as soon as the connection 

1169 is ready. Calling `IOStream` read methods before the socket is 

1170 connected works on some platforms but is non-portable. 

1171 

1172 .. versionchanged:: 4.0 

1173 If no callback is given, returns a `.Future`. 

1174 

1175 .. versionchanged:: 4.2 

1176 SSL certificates are validated by default; pass 

1177 ``ssl_options=dict(cert_reqs=ssl.CERT_NONE)`` or a 

1178 suitably-configured `ssl.SSLContext` to the 

1179 `SSLIOStream` constructor to disable. 

1180 

1181 .. versionchanged:: 6.0 

1182 

1183 The ``callback`` argument was removed. Use the returned 

1184 `.Future` instead. 

1185 

1186 """ 

1187 self._connecting = True 

1188 future = Future() # type: Future[_IOStreamType] 

1189 self._connect_future = typing.cast("Future[IOStream]", future) 

1190 try: 

1191 self.socket.connect(address) 

1192 except BlockingIOError: 

1193 # In non-blocking mode we expect connect() to raise an 

1194 # exception with EINPROGRESS or EWOULDBLOCK. 

1195 pass 

1196 except OSError as e: 

1197 # On freebsd, other errors such as ECONNREFUSED may be 

1198 # returned immediately when attempting to connect to 

1199 # localhost, so handle them the same way as an error 

1200 # reported later in _handle_connect. 

1201 if future is None: 

1202 gen_log.warning("Connect error on fd %s: %s", self.socket.fileno(), e) 

1203 self.close(exc_info=e) 

1204 return future 

1205 self._add_io_state(self.io_loop.WRITE) 

1206 return future 

1207 

1208 def start_tls( 

1209 self, 

1210 server_side: bool, 

1211 ssl_options: Optional[Union[Dict[str, Any], ssl.SSLContext]] = None, 

1212 server_hostname: Optional[str] = None, 

1213 ) -> Awaitable["SSLIOStream"]: 

1214 """Convert this `IOStream` to an `SSLIOStream`. 

1215 

1216 This enables protocols that begin in clear-text mode and 

1217 switch to SSL after some initial negotiation (such as the 

1218 ``STARTTLS`` extension to SMTP and IMAP). 

1219 

1220 This method cannot be used if there are outstanding reads 

1221 or writes on the stream, or if there is any data in the 

1222 IOStream's buffer (data in the operating system's socket 

1223 buffer is allowed). This means it must generally be used 

1224 immediately after reading or writing the last clear-text 

1225 data. It can also be used immediately after connecting, 

1226 before any reads or writes. 

1227 

1228 The ``ssl_options`` argument may be either an `ssl.SSLContext` 

1229 object or a dictionary of keyword arguments for the 

1230 `ssl.SSLContext.wrap_socket` function. The ``server_hostname`` argument 

1231 will be used for certificate validation unless disabled 

1232 in the ``ssl_options``. 

1233 

1234 This method returns a `.Future` whose result is the new 

1235 `SSLIOStream`. After this method has been called, 

1236 any other operation on the original stream is undefined. 

1237 

1238 If a close callback is defined on this stream, it will be 

1239 transferred to the new stream. 

1240 

1241 .. versionadded:: 4.0 

1242 

1243 .. versionchanged:: 4.2 

1244 SSL certificates are validated by default; pass 

1245 ``ssl_options=dict(cert_reqs=ssl.CERT_NONE)`` or a 

1246 suitably-configured `ssl.SSLContext` to disable. 

1247 """ 

1248 if ( 

1249 self._read_future 

1250 or self._write_futures 

1251 or self._connect_future 

1252 or self._closed 

1253 or self._read_buffer 

1254 or self._write_buffer 

1255 ): 

1256 raise ValueError("IOStream is not idle; cannot convert to SSL") 

1257 if ssl_options is None: 

1258 if server_side: 

1259 ssl_options = _server_ssl_defaults 

1260 else: 

1261 ssl_options = _client_ssl_defaults 

1262 

1263 socket = self.socket 

1264 self.io_loop.remove_handler(socket) 

1265 self.socket = None # type: ignore 

1266 socket = ssl_wrap_socket( 

1267 socket, 

1268 ssl_options, 

1269 server_hostname=server_hostname, 

1270 server_side=server_side, 

1271 do_handshake_on_connect=False, 

1272 ) 

1273 orig_close_callback = self._close_callback 

1274 self._close_callback = None 

1275 

1276 future = Future() # type: Future[SSLIOStream] 

1277 ssl_stream = SSLIOStream(socket, ssl_options=ssl_options) 

1278 ssl_stream.set_close_callback(orig_close_callback) 

1279 ssl_stream._ssl_connect_future = future 

1280 ssl_stream.max_buffer_size = self.max_buffer_size 

1281 ssl_stream.read_chunk_size = self.read_chunk_size 

1282 return future 

1283 

1284 def _handle_connect(self) -> None: 

1285 try: 

1286 err = self.socket.getsockopt(socket.SOL_SOCKET, socket.SO_ERROR) 

1287 except OSError as e: 

1288 # Hurd doesn't allow SO_ERROR for loopback sockets because all 

1289 # errors for such sockets are reported synchronously. 

1290 if errno_from_exception(e) == errno.ENOPROTOOPT: 

1291 err = 0 

1292 if err != 0: 

1293 self.error = socket.error(err, os.strerror(err)) 

1294 # IOLoop implementations may vary: some of them return 

1295 # an error state before the socket becomes writable, so 

1296 # in that case a connection failure would be handled by the 

1297 # error path in _handle_events instead of here. 

1298 if self._connect_future is None: 

1299 gen_log.warning( 

1300 "Connect error on fd %s: %s", 

1301 self.socket.fileno(), 

1302 errno.errorcode[err], 

1303 ) 

1304 self.close() 

1305 return 

1306 if self._connect_future is not None: 

1307 future = self._connect_future 

1308 self._connect_future = None 

1309 future_set_result_unless_cancelled(future, self) 

1310 self._connecting = False 

1311 

1312 def set_nodelay(self, value: bool) -> None: 

1313 if self.socket is not None and self.socket.family in ( 

1314 socket.AF_INET, 

1315 socket.AF_INET6, 

1316 ): 

1317 try: 

1318 self.socket.setsockopt( 

1319 socket.IPPROTO_TCP, socket.TCP_NODELAY, 1 if value else 0 

1320 ) 

1321 except OSError as e: 

1322 # Sometimes setsockopt will fail if the socket is closed 

1323 # at the wrong time. This can happen with HTTPServer 

1324 # resetting the value to ``False`` between requests. 

1325 if e.errno != errno.EINVAL and not self._is_connreset(e): 

1326 raise 

1327 

1328 

1329class SSLIOStream(IOStream): 

1330 """A utility class to write to and read from a non-blocking SSL socket. 

1331 

1332 If the socket passed to the constructor is already connected, 

1333 it should be wrapped with:: 

1334 

1335 ssl.SSLContext(...).wrap_socket(sock, do_handshake_on_connect=False, **kwargs) 

1336 

1337 before constructing the `SSLIOStream`. Unconnected sockets will be 

1338 wrapped when `IOStream.connect` is finished. 

1339 """ 

1340 

1341 socket = None # type: ssl.SSLSocket 

1342 

1343 def __init__(self, *args: Any, **kwargs: Any) -> None: 

1344 """The ``ssl_options`` keyword argument may either be an 

1345 `ssl.SSLContext` object or a dictionary of keywords arguments 

1346 for `ssl.SSLContext.wrap_socket` 

1347 """ 

1348 self._ssl_options = kwargs.pop("ssl_options", _client_ssl_defaults) 

1349 super().__init__(*args, **kwargs) 

1350 self._ssl_accepting = True 

1351 self._handshake_reading = False 

1352 self._handshake_writing = False 

1353 self._server_hostname = None # type: Optional[str] 

1354 

1355 # If the socket is already connected, attempt to start the handshake. 

1356 try: 

1357 self.socket.getpeername() 

1358 except OSError: 

1359 pass 

1360 else: 

1361 # Indirectly start the handshake, which will run on the next 

1362 # IOLoop iteration and then the real IO state will be set in 

1363 # _handle_events. 

1364 self._add_io_state(self.io_loop.WRITE) 

1365 

1366 def reading(self) -> bool: 

1367 return self._handshake_reading or super().reading() 

1368 

1369 def writing(self) -> bool: 

1370 return self._handshake_writing or super().writing() 

1371 

1372 def _do_ssl_handshake(self) -> None: 

1373 # Based on code from test_ssl.py in the python stdlib 

1374 try: 

1375 self._handshake_reading = False 

1376 self._handshake_writing = False 

1377 self.socket.do_handshake() 

1378 except ssl.SSLError as err: 

1379 if err.args[0] == ssl.SSL_ERROR_WANT_READ: 

1380 self._handshake_reading = True 

1381 return 

1382 elif err.args[0] == ssl.SSL_ERROR_WANT_WRITE: 

1383 self._handshake_writing = True 

1384 return 

1385 elif err.args[0] in (ssl.SSL_ERROR_EOF, ssl.SSL_ERROR_ZERO_RETURN): 

1386 return self.close(exc_info=err) 

1387 elif err.args[0] in (ssl.SSL_ERROR_SSL, ssl.SSL_ERROR_SYSCALL): 

1388 try: 

1389 peer = self.socket.getpeername() 

1390 except Exception: 

1391 peer = "(not connected)" 

1392 gen_log.warning( 

1393 "SSL Error on %s %s: %s", self.socket.fileno(), peer, err 

1394 ) 

1395 return self.close(exc_info=err) 

1396 raise 

1397 except OSError as err: 

1398 # Some port scans (e.g. nmap in -sT mode) have been known 

1399 # to cause do_handshake to raise EBADF and ENOTCONN, so make 

1400 # those errors quiet as well. 

1401 # https://groups.google.com/forum/?fromgroups#!topic/python-tornado/ApucKJat1_0 

1402 # Errno 0 is also possible in some cases (nc -z). 

1403 # https://github.com/tornadoweb/tornado/issues/2504 

1404 if self._is_connreset(err) or err.args[0] in ( 

1405 0, 

1406 errno.EBADF, 

1407 errno.ENOTCONN, 

1408 ): 

1409 return self.close(exc_info=err) 

1410 raise 

1411 except AttributeError as err: 

1412 # On Linux, if the connection was reset before the call to 

1413 # wrap_socket, do_handshake will fail with an 

1414 # AttributeError. 

1415 return self.close(exc_info=err) 

1416 else: 

1417 self._ssl_accepting = False 

1418 # Prior to the introduction of SNI, this is where we would check 

1419 # the server's claimed hostname. 

1420 assert ssl.HAS_SNI 

1421 self._finish_ssl_connect() 

1422 

1423 def _finish_ssl_connect(self) -> None: 

1424 if self._ssl_connect_future is not None: 

1425 future = self._ssl_connect_future 

1426 self._ssl_connect_future = None 

1427 future_set_result_unless_cancelled(future, self) 

1428 

1429 def _handle_read(self) -> None: 

1430 if self._ssl_accepting: 

1431 self._do_ssl_handshake() 

1432 return 

1433 super()._handle_read() 

1434 

1435 def _handle_write(self) -> None: 

1436 if self._ssl_accepting: 

1437 self._do_ssl_handshake() 

1438 return 

1439 super()._handle_write() 

1440 

1441 def connect( 

1442 self, address: Tuple, server_hostname: Optional[str] = None 

1443 ) -> "Future[SSLIOStream]": 

1444 self._server_hostname = server_hostname 

1445 # Ignore the result of connect(). If it fails, 

1446 # wait_for_handshake will raise an error too. This is 

1447 # necessary for the old semantics of the connect callback 

1448 # (which takes no arguments). In 6.0 this can be refactored to 

1449 # be a regular coroutine. 

1450 # TODO: This is trickier than it looks, since if write() 

1451 # is called with a connect() pending, we want the connect 

1452 # to resolve before the write. Or do we care about this? 

1453 # (There's a test for it, but I think in practice users 

1454 # either wait for the connect before performing a write or 

1455 # they don't care about the connect Future at all) 

1456 fut = super().connect(address) 

1457 fut.add_done_callback(lambda f: f.exception()) 

1458 return self.wait_for_handshake() 

1459 

1460 def _handle_connect(self) -> None: 

1461 # Call the superclass method to check for errors. 

1462 super()._handle_connect() 

1463 if self.closed(): 

1464 return 

1465 # When the connection is complete, wrap the socket for SSL 

1466 # traffic. Note that we do this by overriding _handle_connect 

1467 # instead of by passing a callback to super().connect because 

1468 # user callbacks are enqueued asynchronously on the IOLoop, 

1469 # but since _handle_events calls _handle_connect immediately 

1470 # followed by _handle_write we need this to be synchronous. 

1471 # 

1472 # The IOLoop will get confused if we swap out self.socket while the 

1473 # fd is registered, so remove it now and re-register after 

1474 # wrap_socket(). 

1475 self.io_loop.remove_handler(self.socket) 

1476 old_state = self._state 

1477 assert old_state is not None 

1478 self._state = None 

1479 self.socket = ssl_wrap_socket( 

1480 self.socket, 

1481 self._ssl_options, 

1482 server_hostname=self._server_hostname, 

1483 do_handshake_on_connect=False, 

1484 server_side=False, 

1485 ) 

1486 self._add_io_state(old_state) 

1487 

1488 def wait_for_handshake(self) -> "Future[SSLIOStream]": 

1489 """Wait for the initial SSL handshake to complete. 

1490 

1491 If a ``callback`` is given, it will be called with no 

1492 arguments once the handshake is complete; otherwise this 

1493 method returns a `.Future` which will resolve to the 

1494 stream itself after the handshake is complete. 

1495 

1496 Once the handshake is complete, information such as 

1497 the peer's certificate and NPN/ALPN selections may be 

1498 accessed on ``self.socket``. 

1499 

1500 This method is intended for use on server-side streams 

1501 or after using `IOStream.start_tls`; it should not be used 

1502 with `IOStream.connect` (which already waits for the 

1503 handshake to complete). It may only be called once per stream. 

1504 

1505 .. versionadded:: 4.2 

1506 

1507 .. versionchanged:: 6.0 

1508 

1509 The ``callback`` argument was removed. Use the returned 

1510 `.Future` instead. 

1511 

1512 """ 

1513 if self._ssl_connect_future is not None: 

1514 raise RuntimeError("Already waiting") 

1515 future = self._ssl_connect_future = Future() 

1516 if not self._ssl_accepting: 

1517 self._finish_ssl_connect() 

1518 return future 

1519 

1520 def write_to_fd(self, data: memoryview) -> int: 

1521 # clip buffer size at 1GB since SSL sockets only support upto 2GB 

1522 # this change in behaviour is transparent, since the function is 

1523 # already expected to (possibly) write less than the provided buffer 

1524 if len(data) >> 30: 

1525 data = memoryview(data)[: 1 << 30] 

1526 try: 

1527 return self.socket.send(data) # type: ignore 

1528 except ssl.SSLError as e: 

1529 if e.args[0] == ssl.SSL_ERROR_WANT_WRITE: 

1530 # In Python 3.5+, SSLSocket.send raises a WANT_WRITE error if 

1531 # the socket is not writeable; we need to transform this into 

1532 # an EWOULDBLOCK socket.error or a zero return value, 

1533 # either of which will be recognized by the caller of this 

1534 # method. Prior to Python 3.5, an unwriteable socket would 

1535 # simply return 0 bytes written. 

1536 return 0 

1537 raise 

1538 finally: 

1539 # Avoid keeping to data, which can be a memoryview. 

1540 # See https://github.com/tornadoweb/tornado/pull/2008 

1541 del data 

1542 

1543 def read_from_fd(self, buf: Union[bytearray, memoryview]) -> Optional[int]: 

1544 try: 

1545 if self._ssl_accepting: 

1546 # If the handshake hasn't finished yet, there can't be anything 

1547 # to read (attempting to read may or may not raise an exception 

1548 # depending on the SSL version) 

1549 return None 

1550 # clip buffer size at 1GB since SSL sockets only support upto 2GB 

1551 # this change in behaviour is transparent, since the function is 

1552 # already expected to (possibly) read less than the provided buffer 

1553 if len(buf) >> 30: 

1554 buf = memoryview(buf)[: 1 << 30] 

1555 try: 

1556 return self.socket.recv_into(buf, len(buf)) 

1557 except ssl.SSLError as e: 

1558 # SSLError is a subclass of socket.error, so this except 

1559 # block must come first. 

1560 if e.args[0] == ssl.SSL_ERROR_WANT_READ: 

1561 return None 

1562 else: 

1563 raise 

1564 except BlockingIOError: 

1565 return None 

1566 finally: 

1567 del buf 

1568 

1569 def _is_connreset(self, e: BaseException) -> bool: 

1570 if isinstance(e, ssl.SSLError) and e.args[0] == ssl.SSL_ERROR_EOF: 

1571 return True 

1572 return super()._is_connreset(e) 

1573 

1574 

1575class PipeIOStream(BaseIOStream): 

1576 """Pipe-based `IOStream` implementation. 

1577 

1578 The constructor takes an integer file descriptor (such as one returned 

1579 by `os.pipe`) rather than an open file object. Pipes are generally 

1580 one-way, so a `PipeIOStream` can be used for reading or writing but not 

1581 both. 

1582 

1583 ``PipeIOStream`` is only available on Unix-based platforms. 

1584 """ 

1585 

1586 def __init__(self, fd: int, *args: Any, **kwargs: Any) -> None: 

1587 self.fd = fd 

1588 self._fio = io.FileIO(self.fd, "r+") 

1589 if sys.platform == "win32": 

1590 # The form and placement of this assertion is important to mypy. 

1591 # A plain assert statement isn't recognized here. If the assertion 

1592 # were earlier it would worry that the attributes of self aren't 

1593 # set on windows. If it were missing it would complain about 

1594 # the absence of the set_blocking function. 

1595 raise AssertionError("PipeIOStream is not supported on Windows") 

1596 os.set_blocking(fd, False) 

1597 super().__init__(*args, **kwargs) 

1598 

1599 def fileno(self) -> int: 

1600 return self.fd 

1601 

1602 def close_fd(self) -> None: 

1603 self._fio.close() 

1604 

1605 def write_to_fd(self, data: memoryview) -> int: 

1606 try: 

1607 return os.write(self.fd, data) # type: ignore 

1608 finally: 

1609 # Avoid keeping to data, which can be a memoryview. 

1610 # See https://github.com/tornadoweb/tornado/pull/2008 

1611 del data 

1612 

1613 def read_from_fd(self, buf: Union[bytearray, memoryview]) -> Optional[int]: 

1614 try: 

1615 return self._fio.readinto(buf) # type: ignore 

1616 except OSError as e: 

1617 if errno_from_exception(e) == errno.EBADF: 

1618 # If the writing half of a pipe is closed, select will 

1619 # report it as readable but reads will fail with EBADF. 

1620 self.close(exc_info=e) 

1621 return None 

1622 else: 

1623 raise 

1624 finally: 

1625 del buf 

1626 

1627 

1628def doctests() -> Any: 

1629 import doctest 

1630 

1631 return doctest.DocTestSuite()