Module update: Tornade
This commit is contained in:
@@ -96,17 +96,18 @@ class CurlAsyncHTTPClient(AsyncHTTPClient):
|
||||
pycurl.POLL_INOUT: ioloop.IOLoop.READ | ioloop.IOLoop.WRITE
|
||||
}
|
||||
if event == pycurl.POLL_REMOVE:
|
||||
self.io_loop.remove_handler(fd)
|
||||
del self._fds[fd]
|
||||
if fd in self._fds:
|
||||
self.io_loop.remove_handler(fd)
|
||||
del self._fds[fd]
|
||||
else:
|
||||
ioloop_event = event_map[event]
|
||||
if fd not in self._fds:
|
||||
self._fds[fd] = ioloop_event
|
||||
self.io_loop.add_handler(fd, self._handle_events,
|
||||
ioloop_event)
|
||||
else:
|
||||
self._fds[fd] = ioloop_event
|
||||
else:
|
||||
self.io_loop.update_handler(fd, ioloop_event)
|
||||
self._fds[fd] = ioloop_event
|
||||
|
||||
def _set_timeout(self, msecs):
|
||||
"""Called by libcurl to schedule a timeout."""
|
||||
|
||||
+24
-22
@@ -194,7 +194,7 @@ class IOLoop(Configurable):
|
||||
def initialize(self):
|
||||
pass
|
||||
|
||||
def close(self, all_fds=False):
|
||||
def close(self, all_fds = False):
|
||||
"""Closes the IOLoop, freeing any resources used.
|
||||
|
||||
If ``all_fds`` is true, all file descriptors registered on the
|
||||
@@ -320,7 +320,7 @@ class IOLoop(Configurable):
|
||||
"""
|
||||
raise NotImplementedError()
|
||||
|
||||
def add_callback(self, callback):
|
||||
def add_callback(self, callback, *args, **kwargs):
|
||||
"""Calls the given callback on the next I/O loop iteration.
|
||||
|
||||
It is safe to call this method from any thread at any time,
|
||||
@@ -335,7 +335,7 @@ class IOLoop(Configurable):
|
||||
"""
|
||||
raise NotImplementedError()
|
||||
|
||||
def add_callback_from_signal(self, callback):
|
||||
def add_callback_from_signal(self, callback, *args, **kwargs):
|
||||
"""Calls the given callback on the next I/O loop iteration.
|
||||
|
||||
Safe for use from a Python signal handler; should not be used
|
||||
@@ -359,8 +359,7 @@ class IOLoop(Configurable):
|
||||
assert isinstance(future, IOLoop._FUTURE_TYPES)
|
||||
callback = stack_context.wrap(callback)
|
||||
future.add_done_callback(
|
||||
lambda future: self.add_callback(
|
||||
functools.partial(callback, future)))
|
||||
lambda future: self.add_callback(callback, future))
|
||||
|
||||
def _run_callback(self, callback):
|
||||
"""Runs a callback with error handling.
|
||||
@@ -382,7 +381,7 @@ class IOLoop(Configurable):
|
||||
The exception itself is not passed explicitly, but is available
|
||||
in sys.exc_info.
|
||||
"""
|
||||
app_log.error("Exception in callback %r", callback, exc_info=True)
|
||||
app_log.error("Exception in callback %r", callback, exc_info = True)
|
||||
|
||||
|
||||
|
||||
@@ -393,7 +392,7 @@ class PollIOLoop(IOLoop):
|
||||
(Linux), `tornado.platform.kqueue.KQueueIOLoop` (BSD and Mac), or
|
||||
`tornado.platform.select.SelectIOLoop` (all platforms).
|
||||
"""
|
||||
def initialize(self, impl, time_func=None):
|
||||
def initialize(self, impl, time_func = None):
|
||||
super(PollIOLoop, self).initialize()
|
||||
self._impl = impl
|
||||
if hasattr(self._impl, 'fileno'):
|
||||
@@ -417,7 +416,7 @@ class PollIOLoop(IOLoop):
|
||||
lambda fd, events: self._waker.consume(),
|
||||
self.READ)
|
||||
|
||||
def close(self, all_fds=False):
|
||||
def close(self, all_fds = False):
|
||||
with self._callback_lock:
|
||||
self._closing = True
|
||||
self.remove_handler(self._waker.fileno())
|
||||
@@ -426,7 +425,7 @@ class PollIOLoop(IOLoop):
|
||||
try:
|
||||
os.close(fd)
|
||||
except Exception:
|
||||
gen_log.debug("error closing fd %s", fd, exc_info=True)
|
||||
gen_log.debug("error closing fd %s", fd, exc_info = True)
|
||||
self._waker.close()
|
||||
self._impl.close()
|
||||
|
||||
@@ -442,8 +441,8 @@ class PollIOLoop(IOLoop):
|
||||
self._events.pop(fd, None)
|
||||
try:
|
||||
self._impl.unregister(fd)
|
||||
except (OSError, IOError):
|
||||
gen_log.debug("Error deleting fd from IOLoop", exc_info=True)
|
||||
except Exception:
|
||||
gen_log.debug("Error deleting fd from IOLoop", exc_info = True)
|
||||
|
||||
def set_blocking_signal_threshold(self, seconds, action):
|
||||
if not hasattr(signal, "setitimer"):
|
||||
@@ -501,7 +500,7 @@ class PollIOLoop(IOLoop):
|
||||
# IOLoop is just started once at the beginning.
|
||||
signal.set_wakeup_fd(old_wakeup_fd)
|
||||
old_wakeup_fd = None
|
||||
except ValueError: # non-main thread
|
||||
except ValueError: # non-main thread
|
||||
pass
|
||||
|
||||
while True:
|
||||
@@ -569,17 +568,18 @@ class PollIOLoop(IOLoop):
|
||||
while self._events:
|
||||
fd, events = self._events.popitem()
|
||||
try:
|
||||
self._handlers[fd](fd, events)
|
||||
hdlr = self._handlers.get(fd)
|
||||
if hdlr: hdlr(fd, events)
|
||||
except (OSError, IOError), e:
|
||||
if e.args[0] == errno.EPIPE:
|
||||
# Happens when the client closes the connection
|
||||
pass
|
||||
else:
|
||||
app_log.error("Exception in I/O handler for fd %s",
|
||||
fd, exc_info=True)
|
||||
fd, exc_info = True)
|
||||
except Exception:
|
||||
app_log.error("Exception in I/O handler for fd %s",
|
||||
fd, exc_info=True)
|
||||
fd, exc_info = True)
|
||||
# reset the stopped flag so another start/stop pair can be issued
|
||||
self._stopped = False
|
||||
if self._blocking_signal_threshold is not None:
|
||||
@@ -609,12 +609,13 @@ class PollIOLoop(IOLoop):
|
||||
# collection pass whenever there are too many dead timeouts.
|
||||
timeout.callback = None
|
||||
|
||||
def add_callback(self, callback):
|
||||
def add_callback(self, callback, *args, **kwargs):
|
||||
with self._callback_lock:
|
||||
if self._closing:
|
||||
raise RuntimeError("IOLoop is closing")
|
||||
list_empty = not self._callbacks
|
||||
self._callbacks.append(stack_context.wrap(callback))
|
||||
self._callbacks.append(functools.partial(
|
||||
stack_context.wrap(callback), *args, **kwargs))
|
||||
if list_empty and thread.get_ident() != self._thread_ident:
|
||||
# If we're in the IOLoop's thread, we know it's not currently
|
||||
# polling. If we're not, and we added the first callback to an
|
||||
@@ -624,12 +625,12 @@ class PollIOLoop(IOLoop):
|
||||
# avoid it when we can.
|
||||
self._waker.wake()
|
||||
|
||||
def add_callback_from_signal(self, callback):
|
||||
def add_callback_from_signal(self, callback, *args, **kwargs):
|
||||
with stack_context.NullContext():
|
||||
if thread.get_ident() != self._thread_ident:
|
||||
# if the signal is handled on another thread, we can add
|
||||
# it normally (modulo the NullContext)
|
||||
self.add_callback(callback)
|
||||
self.add_callback(callback, *args, **kwargs)
|
||||
else:
|
||||
# If we're on the IOLoop's thread, we cannot use
|
||||
# the regular add_callback because it may deadlock on
|
||||
@@ -639,7 +640,8 @@ class PollIOLoop(IOLoop):
|
||||
# _callback_lock block in IOLoop.start, we may modify
|
||||
# either the old or new version of self._callbacks,
|
||||
# but either way will work.
|
||||
self._callbacks.append(stack_context.wrap(callback))
|
||||
self._callbacks.append(functools.partial(
|
||||
stack_context.wrap(callback), *args, **kwargs))
|
||||
|
||||
|
||||
class _Timeout(object):
|
||||
@@ -682,7 +684,7 @@ class PeriodicCallback(object):
|
||||
|
||||
`start` must be called after the PeriodicCallback is created.
|
||||
"""
|
||||
def __init__(self, callback, callback_time, io_loop=None):
|
||||
def __init__(self, callback, callback_time, io_loop = None):
|
||||
self.callback = callback
|
||||
if callback_time <= 0:
|
||||
raise ValueError("Periodic callback must have a positive callback_time")
|
||||
@@ -710,7 +712,7 @@ class PeriodicCallback(object):
|
||||
try:
|
||||
self.callback()
|
||||
except Exception:
|
||||
app_log.error("Error in periodic callback", exc_info=True)
|
||||
app_log.error("Error in periodic callback", exc_info = True)
|
||||
self._schedule_next()
|
||||
|
||||
def _schedule_next(self):
|
||||
|
||||
+23
-15
@@ -209,11 +209,19 @@ class BaseIOStream(object):
|
||||
"""Call the given callback when the stream is closed."""
|
||||
self._close_callback = stack_context.wrap(callback)
|
||||
|
||||
def close(self):
|
||||
"""Close this stream."""
|
||||
def close(self, exc_info=False):
|
||||
"""Close this stream.
|
||||
|
||||
If ``exc_info`` is true, set the ``error`` attribute to the current
|
||||
exception from `sys.exc_info()` (or if ``exc_info`` is a tuple,
|
||||
use that instead of `sys.exc_info`).
|
||||
"""
|
||||
if not self.closed():
|
||||
if any(sys.exc_info()):
|
||||
self.error = sys.exc_info()[1]
|
||||
if exc_info:
|
||||
if not isinstance(exc_info, tuple):
|
||||
exc_info = sys.exc_info()
|
||||
if any(exc_info):
|
||||
self.error = exc_info[1]
|
||||
if self._read_until_close:
|
||||
callback = self._read_callback
|
||||
self._read_callback = None
|
||||
@@ -285,7 +293,7 @@ class BaseIOStream(object):
|
||||
except Exception:
|
||||
gen_log.error("Uncaught exception, closing connection.",
|
||||
exc_info=True)
|
||||
self.close()
|
||||
self.close(exc_info=True)
|
||||
raise
|
||||
|
||||
def _run_callback(self, callback, *args):
|
||||
@@ -300,7 +308,7 @@ class BaseIOStream(object):
|
||||
# (It would eventually get closed when the socket object is
|
||||
# gc'd, but we don't want to rely on gc happening before we
|
||||
# run out of file descriptors)
|
||||
self.close()
|
||||
self.close(exc_info=True)
|
||||
# Re-raise the exception so that IOLoop.handle_callback_exception
|
||||
# can see it and log the error
|
||||
raise
|
||||
@@ -348,7 +356,7 @@ class BaseIOStream(object):
|
||||
self._pending_callbacks -= 1
|
||||
except Exception:
|
||||
gen_log.warning("error on read", exc_info=True)
|
||||
self.close()
|
||||
self.close(exc_info=True)
|
||||
return
|
||||
if self._read_from_buffer():
|
||||
return
|
||||
@@ -397,9 +405,9 @@ class BaseIOStream(object):
|
||||
# Treat ECONNRESET as a connection close rather than
|
||||
# an error to minimize log spam (the exception will
|
||||
# be available on self.error for apps that care).
|
||||
self.close()
|
||||
self.close(exc_info=True)
|
||||
return
|
||||
self.close()
|
||||
self.close(exc_info=True)
|
||||
raise
|
||||
if chunk is None:
|
||||
return 0
|
||||
@@ -503,7 +511,7 @@ class BaseIOStream(object):
|
||||
else:
|
||||
gen_log.warning("Write error on %d: %s",
|
||||
self.fileno(), e)
|
||||
self.close()
|
||||
self.close(exc_info=True)
|
||||
return
|
||||
if not self._write_buffer and self._write_callback:
|
||||
callback = self._write_callback
|
||||
@@ -664,7 +672,7 @@ class IOStream(BaseIOStream):
|
||||
if e.args[0] not in (errno.EINPROGRESS, errno.EWOULDBLOCK):
|
||||
gen_log.warning("Connect error on fd %d: %s",
|
||||
self.socket.fileno(), e)
|
||||
self.close()
|
||||
self.close(exc_info=True)
|
||||
return
|
||||
self._connect_callback = stack_context.wrap(callback)
|
||||
self._add_io_state(self.io_loop.WRITE)
|
||||
@@ -733,7 +741,7 @@ class SSLIOStream(IOStream):
|
||||
return
|
||||
elif err.args[0] in (ssl.SSL_ERROR_EOF,
|
||||
ssl.SSL_ERROR_ZERO_RETURN):
|
||||
return self.close()
|
||||
return self.close(exc_info=True)
|
||||
elif err.args[0] == ssl.SSL_ERROR_SSL:
|
||||
try:
|
||||
peer = self.socket.getpeername()
|
||||
@@ -741,11 +749,11 @@ class SSLIOStream(IOStream):
|
||||
peer = '(not connected)'
|
||||
gen_log.warning("SSL Error on %d %s: %s",
|
||||
self.socket.fileno(), peer, err)
|
||||
return self.close()
|
||||
return self.close(exc_info=True)
|
||||
raise
|
||||
except socket.error, err:
|
||||
if err.args[0] in (errno.ECONNABORTED, errno.ECONNRESET):
|
||||
return self.close()
|
||||
return self.close(exc_info=True)
|
||||
else:
|
||||
self._ssl_accepting = False
|
||||
if self._ssl_connect_callback is not None:
|
||||
@@ -842,7 +850,7 @@ class PipeIOStream(BaseIOStream):
|
||||
elif e.args[0] == errno.EBADF:
|
||||
# If the writing half of a pipe is closed, select will
|
||||
# report it as readable but reads will fail with EBADF.
|
||||
self.close()
|
||||
self.close(exc_info=True)
|
||||
return None
|
||||
else:
|
||||
raise
|
||||
|
||||
@@ -431,6 +431,8 @@ class TwistedIOLoop(tornado.ioloop.IOLoop):
|
||||
self.reactor.removeWriter(self.fds[fd])
|
||||
|
||||
def remove_handler(self, fd):
|
||||
if fd not in self.fds:
|
||||
return
|
||||
self.fds[fd].lost = True
|
||||
if self.fds[fd].reading:
|
||||
self.reactor.removeReader(self.fds[fd])
|
||||
@@ -444,6 +446,12 @@ class TwistedIOLoop(tornado.ioloop.IOLoop):
|
||||
def stop(self):
|
||||
self.reactor.crash()
|
||||
|
||||
def _run_callback(self, callback, *args, **kwargs):
|
||||
try:
|
||||
callback(*args, **kwargs)
|
||||
except Exception:
|
||||
self.handle_callback_exception(callback)
|
||||
|
||||
def add_timeout(self, deadline, callback):
|
||||
if isinstance(deadline, (int, long, float)):
|
||||
delay = max(deadline - self.time(), 0)
|
||||
@@ -451,13 +459,14 @@ class TwistedIOLoop(tornado.ioloop.IOLoop):
|
||||
delay = deadline.total_seconds()
|
||||
else:
|
||||
raise TypeError("Unsupported deadline %r")
|
||||
return self.reactor.callLater(delay, wrap(callback))
|
||||
return self.reactor.callLater(delay, self._run_callback, wrap(callback))
|
||||
|
||||
def remove_timeout(self, timeout):
|
||||
timeout.cancel()
|
||||
|
||||
def add_callback(self, callback):
|
||||
self.reactor.callFromThread(wrap(callback))
|
||||
def add_callback(self, callback, *args, **kwargs):
|
||||
self.reactor.callFromThread(self._run_callback,
|
||||
wrap(callback), *args, **kwargs)
|
||||
|
||||
def add_callback_from_signal(self, callback):
|
||||
self.add_callback(callback)
|
||||
def add_callback_from_signal(self, callback, *args, **kwargs):
|
||||
self.add_callback(callback, *args, **kwargs)
|
||||
|
||||
@@ -268,7 +268,7 @@ class Subprocess(object):
|
||||
assert ret_pid == pid
|
||||
subproc = cls._waiting.pop(pid)
|
||||
subproc.io_loop.add_callback_from_signal(
|
||||
functools.partial(subproc._set_returncode, status))
|
||||
subproc._set_returncode, status)
|
||||
|
||||
def _set_returncode(self, status):
|
||||
if os.WIFSIGNALED(status):
|
||||
|
||||
@@ -12,7 +12,6 @@ from tornado.util import b, GzipDecompressor
|
||||
|
||||
import base64
|
||||
import collections
|
||||
import contextlib
|
||||
import copy
|
||||
import functools
|
||||
import os.path
|
||||
@@ -134,7 +133,7 @@ class _HTTPConnection(object):
|
||||
self._decompressor = None
|
||||
# Timeout handle returned by IOLoop.add_timeout
|
||||
self._timeout = None
|
||||
with stack_context.StackContext(self.cleanup):
|
||||
with stack_context.ExceptionStackContext(self._handle_exception):
|
||||
self.parsed = urlparse.urlsplit(_unicode(self.request.url))
|
||||
if ssl is None and self.parsed.scheme == "https":
|
||||
raise ValueError("HTTPS requires either python2.6+ or "
|
||||
@@ -309,19 +308,24 @@ class _HTTPConnection(object):
|
||||
if self.final_callback is not None:
|
||||
final_callback = self.final_callback
|
||||
self.final_callback = None
|
||||
final_callback(response)
|
||||
self.io_loop.add_callback(final_callback, response)
|
||||
|
||||
@contextlib.contextmanager
|
||||
def cleanup(self):
|
||||
try:
|
||||
yield
|
||||
except Exception, e:
|
||||
gen_log.warning("uncaught exception", exc_info=True)
|
||||
self._run_callback(HTTPResponse(self.request, 599, error=e,
|
||||
def _handle_exception(self, typ, value, tb):
|
||||
if self.final_callback:
|
||||
gen_log.warning("uncaught exception", exc_info=(typ, value, tb))
|
||||
self._run_callback(HTTPResponse(self.request, 599, error=value,
|
||||
request_time=self.io_loop.time() - self.start_time,
|
||||
))
|
||||
|
||||
if hasattr(self, "stream"):
|
||||
self.stream.close()
|
||||
return True
|
||||
else:
|
||||
# If our callback has already been called, we are probably
|
||||
# catching an exception that is not caused by us but rather
|
||||
# some child of our callback. Rather than drop it on the floor,
|
||||
# pass it along.
|
||||
return False
|
||||
|
||||
def _on_close(self):
|
||||
if self.final_callback is not None:
|
||||
|
||||
+6
-10
@@ -36,9 +36,8 @@ except ImportError:
|
||||
netutil = None
|
||||
SimpleAsyncHTTPClient = None
|
||||
from tornado.log import gen_log
|
||||
from tornado.stack_context import StackContext
|
||||
from tornado.stack_context import ExceptionStackContext
|
||||
from tornado.util import raise_exc_info
|
||||
import contextlib
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
@@ -167,13 +166,10 @@ class AsyncTestCase(unittest.TestCase):
|
||||
'''
|
||||
return IOLoop()
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _stack_context(self):
|
||||
try:
|
||||
yield
|
||||
except Exception:
|
||||
self.__failure = sys.exc_info()
|
||||
self.stop()
|
||||
def _handle_exception(self, typ, value, tb):
|
||||
self.__failure = sys.exc_info()
|
||||
self.stop()
|
||||
return True
|
||||
|
||||
def __rethrow(self):
|
||||
if self.__failure is not None:
|
||||
@@ -182,7 +178,7 @@ class AsyncTestCase(unittest.TestCase):
|
||||
raise_exc_info(failure)
|
||||
|
||||
def run(self, result=None):
|
||||
with StackContext(self._stack_context):
|
||||
with ExceptionStackContext(self._handle_exception):
|
||||
super(AsyncTestCase, self).run(result)
|
||||
# In case an exception escaped super.run or the StackContext caught
|
||||
# an exception when there wasn't a wait() to re-raise it, do so here.
|
||||
|
||||
+7
-8
@@ -1317,10 +1317,8 @@ class Application(object):
|
||||
def add_handlers(self, host_pattern, host_handlers):
|
||||
"""Appends the given handlers to our handler list.
|
||||
|
||||
Note that host patterns are processed sequentially in the
|
||||
order they were added, and only the first matching pattern is
|
||||
used. This means that all handlers for a given host must be
|
||||
added in a single add_handlers call.
|
||||
Host patterns are processed sequentially in the order they were
|
||||
added. All matching patterns will be considered.
|
||||
"""
|
||||
if not host_pattern.endswith("$"):
|
||||
host_pattern += "$"
|
||||
@@ -1365,15 +1363,16 @@ class Application(object):
|
||||
|
||||
def _get_host_handlers(self, request):
|
||||
host = request.host.lower().split(':')[0]
|
||||
matches = []
|
||||
for pattern, handlers in self.handlers:
|
||||
if pattern.match(host):
|
||||
return handlers
|
||||
matches.extend(handlers)
|
||||
# Look for default host if not behind load balancer (for debugging)
|
||||
if "X-Real-Ip" not in request.headers:
|
||||
if not matches and "X-Real-Ip" not in request.headers:
|
||||
for pattern, handlers in self.handlers:
|
||||
if pattern.match(self.default_host):
|
||||
return handlers
|
||||
return None
|
||||
matches.extend(handlers)
|
||||
return matches or None
|
||||
|
||||
def _load_ui_methods(self, methods):
|
||||
if type(methods) is types.ModuleType:
|
||||
|
||||
Reference in New Issue
Block a user