Update tornado
This commit is contained in:
Regular → Executable
+156
-120
@@ -16,11 +16,12 @@
|
||||
|
||||
"""A utility class to write to and read from a non-blocking socket."""
|
||||
|
||||
from __future__ import with_statement
|
||||
from __future__ import absolute_import, division, with_statement
|
||||
|
||||
import collections
|
||||
import errno
|
||||
import logging
|
||||
import os
|
||||
import socket
|
||||
import sys
|
||||
import re
|
||||
@@ -30,16 +31,17 @@ from tornado import stack_context
|
||||
from tornado.util import b, bytes_type
|
||||
|
||||
try:
|
||||
import ssl # Python 2.6+
|
||||
import ssl # Python 2.6+
|
||||
except ImportError:
|
||||
ssl = None
|
||||
|
||||
|
||||
class IOStream(object):
|
||||
r"""A utility class to write to and read from a non-blocking socket.
|
||||
|
||||
We support a non-blocking ``write()`` and a family of ``read_*()`` methods.
|
||||
All of the methods take callbacks (since writing and reading are
|
||||
non-blocking and asynchronous).
|
||||
non-blocking and asynchronous).
|
||||
|
||||
The socket parameter may either be connected or unconnected. For
|
||||
server operations the socket is the result of calling socket.accept().
|
||||
@@ -47,6 +49,9 @@ class IOStream(object):
|
||||
and may either be connected before passing it to the IOStream or
|
||||
connected with IOStream.connect.
|
||||
|
||||
When a stream is closed due to an error, the IOStream's `error`
|
||||
attribute contains the exception object.
|
||||
|
||||
A very simple (and broken) HTTP client using this class::
|
||||
|
||||
from tornado import ioloop
|
||||
@@ -83,6 +88,7 @@ class IOStream(object):
|
||||
self.io_loop = io_loop or ioloop.IOLoop.instance()
|
||||
self.max_buffer_size = max_buffer_size
|
||||
self.read_chunk_size = read_chunk_size
|
||||
self.error = None
|
||||
self._read_buffer = collections.deque()
|
||||
self._write_buffer = collections.deque()
|
||||
self._read_buffer_size = 0
|
||||
@@ -136,31 +142,15 @@ class IOStream(object):
|
||||
|
||||
def read_until_regex(self, regex, callback):
|
||||
"""Call callback when we read the given regex pattern."""
|
||||
assert not self._read_callback, "Already reading"
|
||||
self._set_read_callback(callback)
|
||||
self._read_regex = re.compile(regex)
|
||||
self._read_callback = stack_context.wrap(callback)
|
||||
while True:
|
||||
# See if we've already got the data from a previous read
|
||||
if self._read_from_buffer():
|
||||
return
|
||||
self._check_closed()
|
||||
if self._read_to_buffer() == 0:
|
||||
break
|
||||
self._add_io_state(self.io_loop.READ)
|
||||
|
||||
self._try_inline_read()
|
||||
|
||||
def read_until(self, delimiter, callback):
|
||||
"""Call callback when we read the given delimiter."""
|
||||
assert not self._read_callback, "Already reading"
|
||||
self._set_read_callback(callback)
|
||||
self._read_delimiter = delimiter
|
||||
self._read_callback = stack_context.wrap(callback)
|
||||
while True:
|
||||
# See if we've already got the data from a previous read
|
||||
if self._read_from_buffer():
|
||||
return
|
||||
self._check_closed()
|
||||
if self._read_to_buffer() == 0:
|
||||
break
|
||||
self._add_io_state(self.io_loop.READ)
|
||||
self._try_inline_read()
|
||||
|
||||
def read_bytes(self, num_bytes, callback, streaming_callback=None):
|
||||
"""Call callback when we read the given number of bytes.
|
||||
@@ -169,18 +159,11 @@ class IOStream(object):
|
||||
of data as they become available, and the argument to the final
|
||||
``callback`` will be empty.
|
||||
"""
|
||||
assert not self._read_callback, "Already reading"
|
||||
self._set_read_callback(callback)
|
||||
assert isinstance(num_bytes, (int, long))
|
||||
self._read_bytes = num_bytes
|
||||
self._read_callback = stack_context.wrap(callback)
|
||||
self._streaming_callback = stack_context.wrap(streaming_callback)
|
||||
while True:
|
||||
if self._read_from_buffer():
|
||||
return
|
||||
self._check_closed()
|
||||
if self._read_to_buffer() == 0:
|
||||
break
|
||||
self._add_io_state(self.io_loop.READ)
|
||||
self._try_inline_read()
|
||||
|
||||
def read_until_close(self, callback, streaming_callback=None):
|
||||
"""Reads all data from the socket until it is closed.
|
||||
@@ -192,12 +175,12 @@ class IOStream(object):
|
||||
Subject to ``max_buffer_size`` limit from `IOStream` constructor if
|
||||
a ``streaming_callback`` is not used.
|
||||
"""
|
||||
assert not self._read_callback, "Already reading"
|
||||
self._set_read_callback(callback)
|
||||
if self.closed():
|
||||
self._run_callback(callback, self._consume(self._read_buffer_size))
|
||||
self._read_callback = None
|
||||
return
|
||||
self._read_until_close = True
|
||||
self._read_callback = stack_context.wrap(callback)
|
||||
self._streaming_callback = stack_context.wrap(streaming_callback)
|
||||
self._add_io_state(self.io_loop.READ)
|
||||
|
||||
@@ -211,10 +194,18 @@ class IOStream(object):
|
||||
"""
|
||||
assert isinstance(data, bytes_type)
|
||||
self._check_closed()
|
||||
# We use bool(_write_buffer) as a proxy for write_buffer_size>0,
|
||||
# so never put empty strings in the buffer.
|
||||
if data:
|
||||
# We use bool(_write_buffer) as a proxy for write_buffer_size>0,
|
||||
# so never put empty strings in the buffer.
|
||||
self._write_buffer.append(data)
|
||||
# Break up large contiguous strings before inserting them in the
|
||||
# write buffer, so we don't have to recopy the entire thing
|
||||
# as we slice off pieces to send to the socket.
|
||||
WRITE_BUFFER_CHUNK_SIZE = 128 * 1024
|
||||
if len(data) > WRITE_BUFFER_CHUNK_SIZE:
|
||||
for i in range(0, len(data), WRITE_BUFFER_CHUNK_SIZE):
|
||||
self._write_buffer.append(data[i:i + WRITE_BUFFER_CHUNK_SIZE])
|
||||
else:
|
||||
self._write_buffer.append(data)
|
||||
self._write_callback = stack_context.wrap(callback)
|
||||
self._handle_write()
|
||||
if self._write_buffer:
|
||||
@@ -228,6 +219,8 @@ class IOStream(object):
|
||||
def close(self):
|
||||
"""Close this stream."""
|
||||
if self.socket is not None:
|
||||
if any(sys.exc_info()):
|
||||
self.error = sys.exc_info()[1]
|
||||
if self._read_until_close:
|
||||
callback = self._read_callback
|
||||
self._read_callback = None
|
||||
@@ -239,12 +232,16 @@ class IOStream(object):
|
||||
self._state = None
|
||||
self.socket.close()
|
||||
self.socket = None
|
||||
if self._close_callback and self._pending_callbacks == 0:
|
||||
# if there are pending callbacks, don't run the close callback
|
||||
# until they're done (see _maybe_add_error_handler)
|
||||
cb = self._close_callback
|
||||
self._close_callback = None
|
||||
self._run_callback(cb)
|
||||
self._maybe_run_close_callback()
|
||||
|
||||
def _maybe_run_close_callback(self):
|
||||
if (self.socket is None and self._close_callback and
|
||||
self._pending_callbacks == 0):
|
||||
# if there are pending callbacks, don't run the close callback
|
||||
# until they're done (see _maybe_add_error_handler)
|
||||
cb = self._close_callback
|
||||
self._close_callback = None
|
||||
self._run_callback(cb)
|
||||
|
||||
def reading(self):
|
||||
"""Returns true if we are currently reading from the stream."""
|
||||
@@ -274,6 +271,9 @@ class IOStream(object):
|
||||
if not self.socket:
|
||||
return
|
||||
if events & self.io_loop.ERROR:
|
||||
errno = self.socket.getsockopt(socket.SOL_SOCKET,
|
||||
socket.SO_ERROR)
|
||||
self.error = socket.error(errno, os.strerror(errno))
|
||||
# We may have queued up a user callback in _handle_read or
|
||||
# _handle_write, so don't close the IOStream until those
|
||||
# callbacks have had a chance to run.
|
||||
@@ -332,22 +332,65 @@ class IOStream(object):
|
||||
self.io_loop.add_callback(wrapper)
|
||||
|
||||
def _handle_read(self):
|
||||
while True:
|
||||
try:
|
||||
try:
|
||||
# Read from the socket until we get EWOULDBLOCK or equivalent.
|
||||
# SSL sockets do some internal buffering, and if the data is
|
||||
# sitting in the SSL object's buffer select() and friends
|
||||
# can't see it; the only way to find out if it's there is to
|
||||
# try to read it.
|
||||
result = self._read_to_buffer()
|
||||
except Exception:
|
||||
self.close()
|
||||
return
|
||||
if result == 0:
|
||||
break
|
||||
else:
|
||||
if self._read_from_buffer():
|
||||
return
|
||||
# Pretend to have a pending callback so that an EOF in
|
||||
# _read_to_buffer doesn't trigger an immediate close
|
||||
# callback. At the end of this method we'll either
|
||||
# estabilsh a real pending callback via
|
||||
# _read_from_buffer or run the close callback.
|
||||
#
|
||||
# We need two try statements here so that
|
||||
# pending_callbacks is decremented before the `except`
|
||||
# clause below (which calls `close` and does need to
|
||||
# trigger the callback)
|
||||
self._pending_callbacks += 1
|
||||
while True:
|
||||
# Read from the socket until we get EWOULDBLOCK or equivalent.
|
||||
# SSL sockets do some internal buffering, and if the data is
|
||||
# sitting in the SSL object's buffer select() and friends
|
||||
# can't see it; the only way to find out if it's there is to
|
||||
# try to read it.
|
||||
if self._read_to_buffer() == 0:
|
||||
break
|
||||
finally:
|
||||
self._pending_callbacks -= 1
|
||||
except Exception:
|
||||
logging.warning("error on read", exc_info=True)
|
||||
self.close()
|
||||
return
|
||||
if self._read_from_buffer():
|
||||
return
|
||||
else:
|
||||
self._maybe_run_close_callback()
|
||||
|
||||
def _set_read_callback(self, callback):
|
||||
assert not self._read_callback, "Already reading"
|
||||
self._read_callback = stack_context.wrap(callback)
|
||||
|
||||
def _try_inline_read(self):
|
||||
"""Attempt to complete the current read operation from buffered data.
|
||||
|
||||
If the read can be completed without blocking, schedules the
|
||||
read callback on the next IOLoop iteration; otherwise starts
|
||||
listening for reads on the socket.
|
||||
"""
|
||||
# See if we've already got the data from a previous read
|
||||
if self._read_from_buffer():
|
||||
return
|
||||
self._check_closed()
|
||||
try:
|
||||
# See comments in _handle_read about incrementing _pending_callbacks
|
||||
self._pending_callbacks += 1
|
||||
while True:
|
||||
if self._read_to_buffer() == 0:
|
||||
break
|
||||
self._check_closed()
|
||||
finally:
|
||||
self._pending_callbacks -= 1
|
||||
if self._read_from_buffer():
|
||||
return
|
||||
self._add_io_state(self.io_loop.READ)
|
||||
|
||||
def _read_from_socket(self):
|
||||
"""Attempts to read from the socket.
|
||||
@@ -397,20 +440,21 @@ class IOStream(object):
|
||||
|
||||
Returns True if the read was completed.
|
||||
"""
|
||||
if self._read_bytes is not None:
|
||||
if self._streaming_callback is not None and self._read_buffer_size:
|
||||
bytes_to_consume = min(self._read_bytes, self._read_buffer_size)
|
||||
if self._streaming_callback is not None and self._read_buffer_size:
|
||||
bytes_to_consume = self._read_buffer_size
|
||||
if self._read_bytes is not None:
|
||||
bytes_to_consume = min(self._read_bytes, bytes_to_consume)
|
||||
self._read_bytes -= bytes_to_consume
|
||||
self._run_callback(self._streaming_callback,
|
||||
self._consume(bytes_to_consume))
|
||||
if self._read_buffer_size >= self._read_bytes:
|
||||
num_bytes = self._read_bytes
|
||||
callback = self._read_callback
|
||||
self._read_callback = None
|
||||
self._streaming_callback = None
|
||||
self._read_bytes = None
|
||||
self._run_callback(callback, self._consume(num_bytes))
|
||||
return True
|
||||
self._run_callback(self._streaming_callback,
|
||||
self._consume(bytes_to_consume))
|
||||
if self._read_bytes is not None and self._read_buffer_size >= self._read_bytes:
|
||||
num_bytes = self._read_bytes
|
||||
callback = self._read_callback
|
||||
self._read_callback = None
|
||||
self._streaming_callback = None
|
||||
self._read_bytes = None
|
||||
self._run_callback(callback, self._consume(num_bytes))
|
||||
return True
|
||||
elif self._read_delimiter is not None:
|
||||
# Multi-byte delimiters (e.g. '\r\n') may straddle two
|
||||
# chunks in the read buffer, so we can't easily find them
|
||||
@@ -420,56 +464,41 @@ class IOStream(object):
|
||||
# to be in the first few chunks. Merge the buffer gradually
|
||||
# since large merges are relatively expensive and get undone in
|
||||
# consume().
|
||||
loc = -1
|
||||
if self._read_buffer:
|
||||
loc = self._read_buffer[0].find(self._read_delimiter)
|
||||
while loc == -1 and len(self._read_buffer) > 1:
|
||||
# Grow by doubling, but don't split the second chunk just
|
||||
# because the first one is small.
|
||||
new_len = max(len(self._read_buffer[0]) * 2,
|
||||
(len(self._read_buffer[0]) +
|
||||
len(self._read_buffer[1])))
|
||||
_merge_prefix(self._read_buffer, new_len)
|
||||
loc = self._read_buffer[0].find(self._read_delimiter)
|
||||
if loc != -1:
|
||||
callback = self._read_callback
|
||||
delimiter_len = len(self._read_delimiter)
|
||||
self._read_callback = None
|
||||
self._streaming_callback = None
|
||||
self._read_delimiter = None
|
||||
self._run_callback(callback,
|
||||
self._consume(loc + delimiter_len))
|
||||
return True
|
||||
while True:
|
||||
loc = self._read_buffer[0].find(self._read_delimiter)
|
||||
if loc != -1:
|
||||
callback = self._read_callback
|
||||
delimiter_len = len(self._read_delimiter)
|
||||
self._read_callback = None
|
||||
self._streaming_callback = None
|
||||
self._read_delimiter = None
|
||||
self._run_callback(callback,
|
||||
self._consume(loc + delimiter_len))
|
||||
return True
|
||||
if len(self._read_buffer) == 1:
|
||||
break
|
||||
_double_prefix(self._read_buffer)
|
||||
elif self._read_regex is not None:
|
||||
m = None
|
||||
if self._read_buffer:
|
||||
m = self._read_regex.search(self._read_buffer[0])
|
||||
while m is None and len(self._read_buffer) > 1:
|
||||
# Grow by doubling, but don't split the second chunk just
|
||||
# because the first one is small.
|
||||
new_len = max(len(self._read_buffer[0]) * 2,
|
||||
(len(self._read_buffer[0]) +
|
||||
len(self._read_buffer[1])))
|
||||
_merge_prefix(self._read_buffer, new_len)
|
||||
m = self._read_regex.search(self._read_buffer[0])
|
||||
_merge_prefix(self._read_buffer, sys.maxint)
|
||||
m = self._read_regex.search(self._read_buffer[0])
|
||||
if m:
|
||||
callback = self._read_callback
|
||||
self._read_callback = None
|
||||
self._streaming_callback = None
|
||||
self._read_regex = None
|
||||
self._run_callback(callback, self._consume(m.end()))
|
||||
return True
|
||||
elif 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))
|
||||
while True:
|
||||
m = self._read_regex.search(self._read_buffer[0])
|
||||
if m is not None:
|
||||
callback = self._read_callback
|
||||
self._read_callback = None
|
||||
self._streaming_callback = None
|
||||
self._read_regex = None
|
||||
self._run_callback(callback, self._consume(m.end()))
|
||||
return True
|
||||
if len(self._read_buffer) == 1:
|
||||
break
|
||||
_double_prefix(self._read_buffer)
|
||||
return False
|
||||
|
||||
def _handle_connect(self):
|
||||
err = self.socket.getsockopt(socket.SOL_SOCKET, socket.SO_ERROR)
|
||||
if err != 0:
|
||||
self.error = socket.error(err, os.strerror(err))
|
||||
# IOLoop implementations may vary: some of them return
|
||||
# an error state before the socket becomes writable, so
|
||||
# in that case a connection failure would be handled by the
|
||||
@@ -537,10 +566,7 @@ class IOStream(object):
|
||||
def _maybe_add_error_listener(self):
|
||||
if self._state is None and self._pending_callbacks == 0:
|
||||
if self.socket is None:
|
||||
cb = self._close_callback
|
||||
if cb is not None:
|
||||
self._close_callback = None
|
||||
self._run_callback(cb)
|
||||
self._maybe_run_close_callback()
|
||||
else:
|
||||
self._add_io_state(ioloop.IOLoop.READ)
|
||||
|
||||
@@ -628,7 +654,7 @@ class SSLIOStream(IOStream):
|
||||
return self.close()
|
||||
raise
|
||||
except socket.error, err:
|
||||
if err.args[0] == errno.ECONNABORTED:
|
||||
if err.args[0] in (errno.ECONNABORTED, errno.ECONNRESET):
|
||||
return self.close()
|
||||
else:
|
||||
self._ssl_accepting = False
|
||||
@@ -655,7 +681,6 @@ class SSLIOStream(IOStream):
|
||||
# until we've completed the SSL handshake (so certificates are
|
||||
# available, etc).
|
||||
|
||||
|
||||
def _read_from_socket(self):
|
||||
if self._ssl_accepting:
|
||||
# If the handshake hasn't finished yet, there can't be anything
|
||||
@@ -686,6 +711,16 @@ class SSLIOStream(IOStream):
|
||||
return None
|
||||
return chunk
|
||||
|
||||
|
||||
def _double_prefix(deque):
|
||||
"""Grow by doubling, but don't split the second chunk just because the
|
||||
first one is small.
|
||||
"""
|
||||
new_len = max(len(deque[0]) * 2,
|
||||
(len(deque[0]) + len(deque[1])))
|
||||
_merge_prefix(deque, new_len)
|
||||
|
||||
|
||||
def _merge_prefix(deque, size):
|
||||
"""Replace the first entries in a deque of strings with a single
|
||||
string of up to size bytes.
|
||||
@@ -723,6 +758,7 @@ def _merge_prefix(deque, size):
|
||||
if not deque:
|
||||
deque.appendleft(b(""))
|
||||
|
||||
|
||||
def doctests():
|
||||
import doctest
|
||||
return doctest.DocTestSuite()
|
||||
|
||||
Reference in New Issue
Block a user