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
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
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.
16"""Utility classes to write to and read from non-blocking files and sockets.
18Contents:
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"""
26import asyncio
27import collections
28import errno
29import io
30import numbers
31import os
32import socket
33import ssl
34import sys
35import re
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
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
57if typing.TYPE_CHECKING:
58 from typing import Deque, List, Type # noqa: F401
60_IOStreamType = TypeVar("_IOStreamType", bound="IOStream")
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)
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 )
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
81_WINDOWS = sys.platform.startswith("win")
84class StreamClosedError(IOError):
85 """Exception raised by `IOStream` methods when the stream is closed.
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.
91 The ``real_error`` attribute contains the underlying error that caused
92 the stream to close (if any).
94 .. versionchanged:: 4.3
95 Added the ``real_error`` attribute.
96 """
98 def __init__(self, real_error: Optional[BaseException] = None) -> None:
99 super().__init__("Stream is closed")
100 self.real_error = real_error
103class UnsatisfiableReadError(Exception):
104 """Exception raised when a read cannot be satisfied.
106 Raised by ``read_until`` and ``read_until_regex`` with a ``max_bytes``
107 argument.
108 """
110 pass
113class StreamBufferFullError(Exception):
114 """Exception raised by `IOStream` methods when the buffer is full."""
117class _StreamBuffer:
118 """
119 A specialized buffer that tries to avoid copies when large pieces
120 of data are encountered.
121 """
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
132 def __len__(self) -> int:
133 return self._size
135 # Data above this size will be appended separately instead
136 # of extending an existing bytearray
137 _large_buf_threshold = 2048
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
159 self._size += size
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"")
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]
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
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
203 assert size == 0
204 self._first_pos = pos
207class BaseIOStream:
208 """A utility class to write to and read from a non-blocking file or socket.
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.
217 When a stream is closed due to an error, the IOStream's ``error``
218 attribute contains the exception object.
220 Subclasses must implement `fileno`, `close_fd`, `write_to_fd`,
221 `read_from_fd`, and optionally `get_fd_error`.
223 """
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.
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.
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
282 def fileno(self) -> Union[int, ioloop._Selectable]:
283 """Returns the file descriptor for this stream."""
284 raise NotImplementedError()
286 def close_fd(self) -> None:
287 """Closes the file underlying this stream.
289 ``close_fd`` is called by `BaseIOStream` and should not be called
290 elsewhere; other users should call `close` instead.
291 """
292 raise NotImplementedError()
294 def write_to_fd(self, data: memoryview) -> int:
295 """Attempts to write ``data`` to the underlying file.
297 Returns the number of bytes written.
298 """
299 raise NotImplementedError()
301 def read_from_fd(self, buf: Union[bytearray, memoryview]) -> Optional[int]:
302 """Attempts to read from the underlying file.
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.
309 .. versionchanged:: 5.0
311 Interface redesigned to take a buffer and return a number
312 of bytes instead of a freshly-allocated object.
313 """
314 raise NotImplementedError()
316 def get_fd_error(self) -> Optional[Exception]:
317 """Returns information about any error on the underlying file.
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
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.
331 The result includes the data that matches the regex and anything
332 that came before it.
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.
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.
342 .. versionchanged:: 6.0
344 The ``callback`` argument was removed. Use the returned
345 `.Future` instead.
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
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.
370 The result includes all the data read including the delimiter.
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.
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.
380 .. versionchanged:: 6.0
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
400 def read_bytes(self, num_bytes: int, partial: bool = False) -> Awaitable[bytes]:
401 """Asynchronously read a number of bytes.
403 If ``partial`` is true, data is returned as soon as we have
404 any bytes to return (but never more than ``num_bytes``)
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.
410 .. versionchanged:: 6.0
412 The ``callback`` and ``streaming_callback`` arguments have
413 been removed. Use the returned `.Future` (and
414 ``partial=True`` for ``streaming_callback``) instead.
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
428 def read_into(self, buf: bytearray, partial: bool = False) -> Awaitable[int]:
429 """Asynchronously read a number of bytes.
431 ``buf`` must be a writable buffer into which data will be read.
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.
437 .. versionadded:: 5.0
439 .. versionchanged:: 6.0
441 The ``callback`` argument was removed. Use the returned
442 `.Future` instead.
444 """
445 future = self._start_read()
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)[:]
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
466 try:
467 self._try_inline_read()
468 except:
469 future.add_done_callback(lambda f: f.exception())
470 raise
471 return future
473 def read_until_close(self) -> Awaitable[bytes]:
474 """Asynchronously reads all data from the socket until it is closed.
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.
480 .. versionchanged:: 4.0
481 The callback argument is now optional and a `.Future` will
482 be returned if it is omitted.
484 .. versionchanged:: 6.0
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.
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
503 def write(self, data: Union[bytes, memoryview]) -> "Future[None]":
504 """Asynchronously write the given data to this stream.
506 This method returns a `.Future` that resolves (with a result
507 of ``None``) when the write has been completed.
509 The ``data`` argument may be of type `bytes` or `memoryview`.
511 .. versionchanged:: 4.0
512 Now returns a `.Future` if no callback is given.
514 .. versionchanged:: 4.5
515 Added support for `memoryview` arguments.
517 .. versionchanged:: 6.0
519 The ``callback`` argument was removed. Use the returned
520 `.Future` instead.
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
545 def set_close_callback(self, callback: Optional[Callable[[], None]]) -> None:
546 """Call the given callback when the stream is closed.
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.
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()
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.
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()
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
658 def reading(self) -> bool:
659 """Returns ``True`` if we are currently reading from the stream."""
660 return self._read_future is not None
662 def writing(self) -> bool:
663 """Returns ``True`` if we are currently writing to the stream."""
664 return bool(self._write_buffer)
666 def closed(self) -> bool:
667 """Returns ``True`` if the stream has been closed."""
668 return self._closed
670 def set_nodelay(self, value: bool) -> None:
671 """Sets the no-delay flag for this stream.
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.
679 This flag is currently defined only for TCP-based ``IOStreams``.
681 .. versionadded:: 3.1
682 """
683 pass
685 def _handle_connect(self) -> None:
686 raise NotImplementedError()
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
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
764 # If we've read all the bytes we can use, break out of
765 # this loop.
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
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()
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)
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
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()
834 def _try_inline_read(self) -> None:
835 """Attempt to complete the current read operation from buffered data.
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)
856 def _read_to_buffer(self) -> Optional[int]:
857 """Reads from the socket and appends the result to the read buffer.
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
905 def _read_from_buffer(self, pos: int) -> None:
906 """Attempts to complete the currently-pending read from the buffer.
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)
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.
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
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 )
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
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
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)
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
1010 def _check_closed(self) -> None:
1011 if self.closed():
1012 raise StreamClosedError(real_error=self.error)
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)
1029 def _add_io_state(self, state: int) -> None:
1030 """Adds `state` (IOLoop.{READ,WRITE} flags) to our event handler.
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.
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.
1045 TODO: reevaluate this now that callbacks are gone.
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)
1058 def _is_connreset(self, exc: BaseException) -> bool:
1059 """Return ``True`` if exc is ECONNRESET or equivalent.
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 )
1069class IOStream(BaseIOStream):
1070 r"""Socket-based `IOStream` implementation.
1072 This class supports the read and write methods from `BaseIOStream`
1073 plus a `connect` method.
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`.
1082 A very simple (and broken) HTTP client using this class:
1084 .. testcode::
1086 import socket
1087 import tornado
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()
1104 if __name__ == '__main__':
1105 asyncio.run(main())
1107 """
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)
1114 def fileno(self) -> Union[int, ioloop._Selectable]:
1115 return self.socket
1117 def close_fd(self) -> None:
1118 self.socket.close()
1119 self.socket = None # type: ignore
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))
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
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
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.
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.
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).
1162 In SSL mode, the ``server_hostname`` parameter will be used
1163 for certificate validation (unless disabled in the
1164 ``ssl_options``) and SNI.
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.
1172 .. versionchanged:: 4.0
1173 If no callback is given, returns a `.Future`.
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.
1181 .. versionchanged:: 6.0
1183 The ``callback`` argument was removed. Use the returned
1184 `.Future` instead.
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
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`.
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).
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.
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``.
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.
1238 If a close callback is defined on this stream, it will be
1239 transferred to the new stream.
1241 .. versionadded:: 4.0
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
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
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
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
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
1329class SSLIOStream(IOStream):
1330 """A utility class to write to and read from a non-blocking SSL socket.
1332 If the socket passed to the constructor is already connected,
1333 it should be wrapped with::
1335 ssl.SSLContext(...).wrap_socket(sock, do_handshake_on_connect=False, **kwargs)
1337 before constructing the `SSLIOStream`. Unconnected sockets will be
1338 wrapped when `IOStream.connect` is finished.
1339 """
1341 socket = None # type: ssl.SSLSocket
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]
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)
1366 def reading(self) -> bool:
1367 return self._handshake_reading or super().reading()
1369 def writing(self) -> bool:
1370 return self._handshake_writing or super().writing()
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()
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)
1429 def _handle_read(self) -> None:
1430 if self._ssl_accepting:
1431 self._do_ssl_handshake()
1432 return
1433 super()._handle_read()
1435 def _handle_write(self) -> None:
1436 if self._ssl_accepting:
1437 self._do_ssl_handshake()
1438 return
1439 super()._handle_write()
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()
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)
1488 def wait_for_handshake(self) -> "Future[SSLIOStream]":
1489 """Wait for the initial SSL handshake to complete.
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.
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``.
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.
1505 .. versionadded:: 4.2
1507 .. versionchanged:: 6.0
1509 The ``callback`` argument was removed. Use the returned
1510 `.Future` instead.
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
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
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
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)
1575class PipeIOStream(BaseIOStream):
1576 """Pipe-based `IOStream` implementation.
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.
1583 ``PipeIOStream`` is only available on Unix-based platforms.
1584 """
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)
1599 def fileno(self) -> int:
1600 return self.fd
1602 def close_fd(self) -> None:
1603 self._fio.close()
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
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
1628def doctests() -> Any:
1629 import doctest
1631 return doctest.DocTestSuite()