Update Tornado
This commit is contained in:
+68
-13
@@ -48,6 +48,12 @@ except ImportError:
|
||||
|
||||
|
||||
class StreamClosedError(IOError):
|
||||
"""Exception raised by `IOStream` methods when the stream is closed.
|
||||
|
||||
Note that the close callback is scheduled to run *after* other
|
||||
callbacks on the stream (to allow for buffered data to be processed),
|
||||
so you may see this error before you see the close callback.
|
||||
"""
|
||||
pass
|
||||
|
||||
|
||||
@@ -64,10 +70,10 @@ class BaseIOStream(object):
|
||||
Subclasses must implement `fileno`, `close_fd`, `write_to_fd`,
|
||||
`read_from_fd`, and optionally `get_fd_error`.
|
||||
"""
|
||||
def __init__(self, io_loop=None, max_buffer_size=104857600,
|
||||
def __init__(self, io_loop=None, max_buffer_size=None,
|
||||
read_chunk_size=4096):
|
||||
self.io_loop = io_loop or ioloop.IOLoop.current()
|
||||
self.max_buffer_size = max_buffer_size
|
||||
self.max_buffer_size = max_buffer_size or 104857600
|
||||
self.read_chunk_size = read_chunk_size
|
||||
self.error = None
|
||||
self._read_buffer = collections.deque()
|
||||
@@ -234,6 +240,10 @@ class BaseIOStream(object):
|
||||
if any(exc_info):
|
||||
self.error = exc_info[1]
|
||||
if self._read_until_close:
|
||||
if (self._streaming_callback is not None and
|
||||
self._read_buffer_size):
|
||||
self._run_callback(self._streaming_callback,
|
||||
self._consume(self._read_buffer_size))
|
||||
callback = self._read_callback
|
||||
self._read_callback = None
|
||||
self._read_until_close = False
|
||||
@@ -269,6 +279,21 @@ class BaseIOStream(object):
|
||||
"""Returns true if the stream has been closed."""
|
||||
return self._closed
|
||||
|
||||
def set_nodelay(self, value):
|
||||
"""Sets the no-delay flag for this stream.
|
||||
|
||||
By default, data written to TCP streams may be held for a time
|
||||
to make the most efficient use of bandwidth (according to
|
||||
Nagle's algorithm). The no-delay flag requests that data be
|
||||
written as soon as possible, even if doing so would consume
|
||||
additional bandwidth.
|
||||
|
||||
This flag is currently defined only for TCP-based ``IOStreams``.
|
||||
|
||||
.. versionadded:: 3.1
|
||||
"""
|
||||
pass
|
||||
|
||||
def _handle_events(self, fd, events):
|
||||
if self.closed():
|
||||
gen_log.warning("Got events for closed stream %d", fd)
|
||||
@@ -392,13 +417,21 @@ class BaseIOStream(object):
|
||||
return
|
||||
self._check_closed()
|
||||
try:
|
||||
# See comments in _handle_read about incrementing _pending_callbacks
|
||||
self._pending_callbacks += 1
|
||||
while not self.closed():
|
||||
if self._read_to_buffer() == 0:
|
||||
break
|
||||
finally:
|
||||
self._pending_callbacks -= 1
|
||||
try:
|
||||
# See comments in _handle_read about incrementing _pending_callbacks
|
||||
self._pending_callbacks += 1
|
||||
while not self.closed():
|
||||
if self._read_to_buffer() == 0:
|
||||
break
|
||||
finally:
|
||||
self._pending_callbacks -= 1
|
||||
except Exception:
|
||||
# If there was an in _read_to_buffer, we called close() already,
|
||||
# but couldn't run the close callback because of _pending_callbacks.
|
||||
# Before we escape from this function, run the close callback if
|
||||
# applicable.
|
||||
self._maybe_run_close_callback()
|
||||
raise
|
||||
if self._read_from_buffer():
|
||||
return
|
||||
self._maybe_add_error_listener()
|
||||
@@ -522,8 +555,12 @@ class BaseIOStream(object):
|
||||
self._write_buffer_frozen = True
|
||||
break
|
||||
else:
|
||||
gen_log.warning("Write error on %d: %s",
|
||||
self.fileno(), e)
|
||||
if e.args[0] not in (errno.EPIPE, errno.ECONNRESET):
|
||||
# Broken pipe errors are usually caused by connection
|
||||
# reset, and its better to not log EPIPE errors to
|
||||
# minimize log spam
|
||||
gen_log.warning("Write error on %d: %s",
|
||||
self.fileno(), e)
|
||||
self.close(exc_info=True)
|
||||
return
|
||||
if not self._write_buffer and self._write_callback:
|
||||
@@ -714,6 +751,19 @@ class IOStream(BaseIOStream):
|
||||
self._run_callback(callback)
|
||||
self._connecting = False
|
||||
|
||||
def set_nodelay(self, value):
|
||||
if (self.socket is not None and
|
||||
self.socket.family in (socket.AF_INET, socket.AF_INET6)):
|
||||
try:
|
||||
self.socket.setsockopt(socket.IPPROTO_TCP,
|
||||
socket.TCP_NODELAY, 1 if value else 0)
|
||||
except socket.error as e:
|
||||
# Sometimes setsockopt will fail if the socket is closed
|
||||
# at the wrong time. This can happen with HTTPServer
|
||||
# resetting the value to false between requests.
|
||||
if e.errno != errno.EINVAL:
|
||||
raise
|
||||
|
||||
|
||||
class SSLIOStream(IOStream):
|
||||
"""A utility class to write to and read from a non-blocking SSL socket.
|
||||
@@ -764,7 +814,7 @@ class SSLIOStream(IOStream):
|
||||
elif err.args[0] == ssl.SSL_ERROR_SSL:
|
||||
try:
|
||||
peer = self.socket.getpeername()
|
||||
except:
|
||||
except Exception:
|
||||
peer = '(not connected)'
|
||||
gen_log.warning("SSL Error on %d %s: %s",
|
||||
self.socket.fileno(), peer, err)
|
||||
@@ -773,6 +823,11 @@ class SSLIOStream(IOStream):
|
||||
except socket.error as err:
|
||||
if err.args[0] in (errno.ECONNABORTED, errno.ECONNRESET):
|
||||
return self.close(exc_info=True)
|
||||
except AttributeError:
|
||||
# On Linux, if the connection was reset before the call to
|
||||
# wrap_socket, do_handshake will fail with an
|
||||
# AttributeError.
|
||||
return self.close(exc_info=True)
|
||||
else:
|
||||
self._ssl_accepting = False
|
||||
if not self._verify_cert(self.socket.getpeercert()):
|
||||
@@ -825,7 +880,7 @@ class SSLIOStream(IOStream):
|
||||
def connect(self, address, callback=None, server_hostname=None):
|
||||
# Save the user's callback and run it after the ssl handshake
|
||||
# has completed.
|
||||
self._ssl_connect_callback = callback
|
||||
self._ssl_connect_callback = stack_context.wrap(callback)
|
||||
self._server_hostname = server_hostname
|
||||
super(SSLIOStream, self).connect(address, callback=None)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user