From 74e4b015a954712323df6c736a13ba713db19b62 Mon Sep 17 00:00:00 2001 From: Ruud Date: Sun, 30 Dec 2012 18:38:52 +0100 Subject: [PATCH] Module update: Tornade --- libs/tornado/curl_httpclient.py | 9 +++--- libs/tornado/ioloop.py | 46 ++++++++++++++++--------------- libs/tornado/iostream.py | 38 +++++++++++++++---------- libs/tornado/platform/twisted.py | 19 +++++++++---- libs/tornado/process.py | 2 +- libs/tornado/simple_httpclient.py | 24 +++++++++------- libs/tornado/testing.py | 16 ++++------- libs/tornado/web.py | 15 +++++----- 8 files changed, 94 insertions(+), 75 deletions(-) diff --git a/libs/tornado/curl_httpclient.py b/libs/tornado/curl_httpclient.py index a6c0bb0d..52350d24 100755 --- a/libs/tornado/curl_httpclient.py +++ b/libs/tornado/curl_httpclient.py @@ -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.""" diff --git a/libs/tornado/ioloop.py b/libs/tornado/ioloop.py index 3c9b05b9..7b320e59 100755 --- a/libs/tornado/ioloop.py +++ b/libs/tornado/ioloop.py @@ -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): diff --git a/libs/tornado/iostream.py b/libs/tornado/iostream.py index 40ac4964..6eec2a35 100755 --- a/libs/tornado/iostream.py +++ b/libs/tornado/iostream.py @@ -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 diff --git a/libs/tornado/platform/twisted.py b/libs/tornado/platform/twisted.py index 6c3cbf96..1efc82b7 100755 --- a/libs/tornado/platform/twisted.py +++ b/libs/tornado/platform/twisted.py @@ -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) diff --git a/libs/tornado/process.py b/libs/tornado/process.py index 9e048c19..fa0be555 100755 --- a/libs/tornado/process.py +++ b/libs/tornado/process.py @@ -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): diff --git a/libs/tornado/simple_httpclient.py b/libs/tornado/simple_httpclient.py index faff83c8..7000d987 100755 --- a/libs/tornado/simple_httpclient.py +++ b/libs/tornado/simple_httpclient.py @@ -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: diff --git a/libs/tornado/testing.py b/libs/tornado/testing.py index 59456433..22376627 100755 --- a/libs/tornado/testing.py +++ b/libs/tornado/testing.py @@ -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. diff --git a/libs/tornado/web.py b/libs/tornado/web.py index 7d45ce53..41ce95d4 100755 --- a/libs/tornado/web.py +++ b/libs/tornado/web.py @@ -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: