From 3e3a66620b278f7c03e36ebedb3591d0aaa7e8f4 Mon Sep 17 00:00:00 2001 From: Pablo Hoffman Date: Sat, 14 Aug 2010 21:10:37 -0300 Subject: [PATCH] Added support for returning deferreds from (some) signal handlers. Closes #193 --- docs/faq.rst | 10 +++- docs/topics/signals.rst | 58 ++++++++++++++++----- scrapy/core/engine.py | 15 +++--- scrapy/core/scraper.py | 16 +++--- scrapy/tests/test_utils_signal.py | 85 ++++++++++++++++++++++--------- scrapy/utils/signal.py | 28 ++++++++++ 6 files changed, 158 insertions(+), 54 deletions(-) diff --git a/docs/faq.rst b/docs/faq.rst index 57fde2deb..9276c312a 100644 --- a/docs/faq.rst +++ b/docs/faq.rst @@ -29,8 +29,8 @@ application framework). Does Scrapy work with Python 3.0? --------------------------------- -No, and there are no plans to port Scrapy to Python 3.0 yet. At the moment -Scrapy works with Python 2.5 or 2.6. +No, and there are no plans to port Scrapy to Python 3.0 yet. At the moment, +Scrapy works with Python 2.5, 2.6 and 2.7. Did Scrapy "steal" X from Django? --------------------------------- @@ -167,3 +167,9 @@ Can I use JSON for large exports? It'll depend on how large your output is. See :ref:`this warning ` in :class:`~scrapy.contrib.exporter.JsonItemExporter` documentation. + +Can I return (Twisted) deferreds from signal handlers? +------------------------------------------------------ + +Some signals support returning deferreds form their handlers, others don't. See +the :ref:`topics-signals-ref` to know which ones. diff --git a/docs/topics/signals.rst b/docs/topics/signals.rst index 482bdf737..25fc100f1 100644 --- a/docs/topics/signals.rst +++ b/docs/topics/signals.rst @@ -4,19 +4,28 @@ Signals ======= -Scrapy uses signals extensively to notify when certain actions occur. You can -catch some of those signals in your Scrapy project or extension to perform -additional tasks or extend Scrapy to add functionality not provided out of the -box. +Scrapy uses signals extensively to notify when certain events occur. You can +catch some of those signals in your Scrapy project (using an :ref:`extension +`, for example) to perform additional tasks or extend Scrapy +to add functionality not provided out of the box. -Even though signals provide several arguments, the handlers which catch them -don't have to receive all of them. +Even though signals provide several arguments, the handlers that catch them +don't need to accept all of them - the signal dispatching mechanism will only +deliver the arguments that the handler receives. -For more information about working when see the documentation of -`pydispatcher`_ (library used to implement signals). +Finally, for more detailed information about signals internals see the +documentation of `pydispatcher`_ (the which the signal dispatching mechanism is +based on). .. _pydispatcher: http://pydispatcher.sourceforge.net/ +Deferred signal handlers +======================== + +Some signals support returning `Twisted deferreds`_ from their handlers, see +the :ref:`topics-signals-ref` below to know which ones. + +.. _Twisted deferreds: http://twistedmatrix.com/documents/current/core/howto/defer.html .. _topics-signals-ref: @@ -26,8 +35,7 @@ Built-in signals reference .. module:: scrapy.signals :synopsis: Signals definitions -Here's a list of signals used in Scrapy and their meaning, in alphabetical -order. +Here's the list of Scrapy built-in signals and their meaning. engine_started -------------- @@ -38,6 +46,8 @@ engine_started Sent when the Scrapy engine is started (for example, when a crawling process has started). + This signal supports returning deferreds from their handlers. + engine_stopped -------------- @@ -47,6 +57,8 @@ engine_stopped Sent when the Scrapy engine is stopped (for example, when a crawling process has finished). + This signal supports returning deferreds from their handlers. + item_scraped ------------ @@ -56,6 +68,8 @@ item_scraped Sent when the engine receives a new scraped item from the spider, and right before the item is sent to the :ref:`topics-item-pipeline`. + This signal supports returning deferreds from their handlers. + :param item: is the item scraped :type item: :class:`~scrapy.item.Item` object @@ -74,6 +88,8 @@ item_passed Sent after an item has passed all the :ref:`topics-item-pipeline` stages without being dropped. + This signal supports returning deferreds from their handlers. + :param item: the item which passed the pipeline :type item: :class:`~scrapy.item.Item` object @@ -93,6 +109,8 @@ item_dropped Sent after an item has been dropped from the :ref:`topics-item-pipeline` when some stage raised a :exc:`~scrapy.exceptions.DropItem` exception. + This signal supports returning deferreds from their handlers. + :param item: the item dropped from the :ref:`topics-item-pipeline` :type item: :class:`~scrapy.item.Item` object @@ -113,6 +131,8 @@ spider_closed Sent after a spider has been closed. This can be used to release per-spider resources reserved on :signal:`spider_opened`. + This signal supports returning deferreds from their handlers. + :param spider: the spider which has been closed :type spider: :class:`~scrapy.spider.BaseSpider` object @@ -135,6 +155,8 @@ spider_opened reserve per-spider resources, but can be used for any task that needs to be performed when a spider is opened. + This signal supports returning deferreds from their handlers. + :param spider: the spider which has been opened :type spider: :class:`~scrapy.spider.BaseSpider` object @@ -157,6 +179,8 @@ spider_idle You can, for example, schedule some requests in your :signal:`spider_idle` handler to prevent the spider from being closed. + This signal does not support returning deferreds from their handlers. + :param spider: the spider which has gone idle :type spider: :class:`~scrapy.spider.BaseSpider` object @@ -168,6 +192,8 @@ request_received Sent when the engine receives a :class:`~scrapy.http.Request` from a spider. + This signal does not support returning deferreds from their handlers. + :param request: the request received :type request: :class:`~scrapy.http.Request` object @@ -182,6 +208,8 @@ request_uploaded Sent right after the download has sent a :class:`~scrapy.http.Request`. + This signal does not support returning deferreds from their handlers. + :param request: the request uploaded/sent :type request: :class:`~scrapy.http.Request` object @@ -194,15 +222,17 @@ response_received .. signal:: response_received .. function:: response_received(response, spider) + Sent when the engine receives a new :class:`~scrapy.http.Response` from the + downloader. + + This signal does not support returning deferreds from their handlers. + :param response: the response received :type response: :class:`~scrapy.http.Response` object :param spider: the spider for which the response is intended :type spider: :class:`~scrapy.spider.BaseSpider` object - Sent when the engine receives a new :class:`~scrapy.http.Response` from the - downloader. - response_downloaded ------------------- @@ -211,6 +241,8 @@ response_downloaded Sent by the downloader right after a ``HTTPResponse`` is downloaded. + This signal does not support returning deferreds from their handlers. + :param response: the response downloaded :type response: :class:`~scrapy.http.Response` object diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 78c6457ee..35163641f 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -19,7 +19,7 @@ from scrapy.exceptions import IgnoreRequest, DontCloseSpider from scrapy.http import Response, Request from scrapy.spider import spiders from scrapy.utils.misc import load_object -from scrapy.utils.signal import send_catch_log +from scrapy.utils.signal import send_catch_log, send_catch_log_deferred from scrapy.utils.defer import mustbe_deferred class ExecutionEngine(object): @@ -44,11 +44,12 @@ class ExecutionEngine(object): self.configured = True self._spider_closed_callback = spider_closed_callback + @defer.inlineCallbacks def start(self): """Start the execution engine""" assert not self.running, "Engine already running" self.start_time = time() - send_catch_log(signal=signals.engine_started) + yield send_catch_log_deferred(signal=signals.engine_started) self.running = True def stop(self): @@ -218,7 +219,7 @@ class ExecutionEngine(object): self.downloader.open_spider(spider) yield self.scraper.open_spider(spider) stats.open_spider(spider) - send_catch_log(signals.spider_opened, spider=spider) + yield send_catch_log_deferred(signals.spider_opened, spider=spider) self.next_request(spider) def _spider_idle(self, spider): @@ -268,11 +269,12 @@ class ExecutionEngine(object): def _finish_closing_spider(self, spider): """This function is called after the spider has been closed""" reason = self.closing.pop(spider, 'finished') - send_catch_log(signal=signals.spider_closed, spider=spider, reason=reason) call = self._next_request_calls.pop(spider, None) if call and call.active(): call.cancel() - dfd = defer.maybeDeferred(stats.close_spider, spider, reason=reason) + dfd = send_catch_log_deferred(signal=signals.spider_closed, \ + spider=spider, reason=reason) + dfd.addBoth(lambda _: stats.close_spider(spider, reason=reason)) dfd.addErrback(log.err, "Unhandled error in stats.close_spider()", spider=spider) dfd.addBoth(lambda _: spiders.close_spider(spider)) @@ -283,5 +285,6 @@ class ExecutionEngine(object): dfd.addBoth(lambda _: self._spider_closed_callback(spider)) return dfd + @defer.inlineCallbacks def _finish_stopping_engine(self): - send_catch_log(signal=signals.engine_stopped) + yield send_catch_log_deferred(signal=signals.engine_stopped) diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index 48d8329b8..1fd5fc05a 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -7,7 +7,7 @@ from twisted.internet import defer from scrapy.utils.defer import defer_result, defer_succeed, parallel, iter_errback from scrapy.utils.spider import iterate_spider_output from scrapy.utils.misc import load_object -from scrapy.utils.signal import send_catch_log +from scrapy.utils.signal import send_catch_log, send_catch_log_deferred from scrapy.exceptions import IgnoreRequest, DropItem from scrapy import signals from scrapy.http import Request, Response @@ -171,14 +171,10 @@ class Scraper(object): elif isinstance(output, BaseItem): log.msg("Scraped %s in <%s>" % (output, request.url), level=log.DEBUG, \ spider=spider) - send_catch_log(signal=signals.item_scraped, \ - item=output, spider=spider, response=response) self.sites[spider].itemproc_size += 1 - # FIXME: this can't be called here because the stats spider may be - # already closed - #stats.max_value('scraper/max_itemproc_size', \ - # self.sites[spider].itemproc_size, spider=spider) - dfd = self.itemproc.process_item(output, spider) + dfd = send_catch_log_deferred(signal=signals.item_scraped, \ + item=output, spider=spider, response=response) + dfd.addBoth(lambda _: self.itemproc.process_item(output, spider)) dfd.addBoth(self._itemproc_finished, output, spider) return dfd elif output is None: @@ -210,12 +206,12 @@ class Scraper(object): ex = output.value if isinstance(ex, DropItem): log.msg("Dropped %s - %s" % (item, str(ex)), level=log.WARNING, spider=spider) - send_catch_log(signal=signals.item_dropped, \ + return send_catch_log_deferred(signal=signals.item_dropped, \ item=item, spider=spider, exception=output.value) else: log.err(output, 'Error processing %s' % item, spider=spider) else: log.msg("Passed %s" % item, log.INFO, spider=spider) - send_catch_log(signal=signals.item_passed, \ + return send_catch_log_deferred(signal=signals.item_passed, \ item=item, spider=spider, output=output) diff --git a/scrapy/tests/test_utils_signal.py b/scrapy/tests/test_utils_signal.py index 640174c57..92f63b874 100644 --- a/scrapy/tests/test_utils_signal.py +++ b/scrapy/tests/test_utils_signal.py @@ -1,45 +1,84 @@ from twisted.trial import unittest from twisted.python import log as txlog from twisted.python.failure import Failure +from twisted.internet import defer, reactor from scrapy.xlib.pydispatch import dispatcher -from scrapy.utils.signal import send_catch_log +from scrapy.utils.signal import send_catch_log, send_catch_log_deferred from scrapy import log -test_signal = object() - -class SignalUtilsTest(unittest.TestCase): +class SendCatchLogTest(unittest.TestCase): + @defer.inlineCallbacks def test_send_catch_log(self): + test_signal = object() handlers_called = set() - def test_handler_error(arg): - handlers_called.add(test_handler_error) - a = 1/0 - - def test_handler_check(arg): - handlers_called.add(test_handler_check) - assert arg == 'test' - return "OK" - def log_received(event): handlers_called.add(log_received) - assert "test_handler_error" in event['message'][0] + assert "error_handler" in event['message'][0] assert event['logLevel'] == log.ERROR txlog.addObserver(log_received) - dispatcher.connect(test_handler_error, signal=test_signal) - dispatcher.connect(test_handler_check, signal=test_signal) - result = send_catch_log(test_signal, arg='test') + dispatcher.connect(self.error_handler, signal=test_signal) + dispatcher.connect(self.ok_handler, signal=test_signal) + result = yield defer.maybeDeferred(self._get_result, test_signal, arg='test', \ + handlers_called=handlers_called) - assert test_handler_error in handlers_called - assert test_handler_check in handlers_called + assert self.error_handler in handlers_called + assert self.ok_handler in handlers_called assert log_received in handlers_called - self.assertEqual(result[0][0], test_handler_error) + self.assertEqual(result[0][0], self.error_handler) self.assert_(isinstance(result[0][1], Failure)) - self.assertEqual(result[1], (test_handler_check, "OK")) + self.assertEqual(result[1], (self.ok_handler, "OK")) txlog.removeObserver(log_received) self.flushLoggedErrors() - dispatcher.disconnect(test_handler_error, signal=test_signal) - dispatcher.disconnect(test_handler_check, signal=test_signal) + dispatcher.disconnect(self.error_handler, signal=test_signal) + dispatcher.disconnect(self.ok_handler, signal=test_signal) + + def _get_result(self, signal, *a, **kw): + return send_catch_log(signal, *a, **kw) + + def error_handler(self, arg, handlers_called): + handlers_called.add(self.error_handler) + a = 1/0 + + def ok_handler(self, arg, handlers_called): + handlers_called.add(self.ok_handler) + assert arg == 'test' + return "OK" + + +class SendCatchLogDeferredTest(SendCatchLogTest): + + def _get_result(self, signal, *a, **kw): + return send_catch_log_deferred(signal, *a, **kw) + + +class SendCatchLogDeferredTest2(SendCatchLogTest): + + def ok_handler(self, arg, handlers_called): + handlers_called.add(self.ok_handler) + assert arg == 'test' + d = defer.Deferred() + reactor.callLater(0, d.callback, "OK") + return d + + def _get_result(self, signal, *a, **kw): + return send_catch_log_deferred(signal, *a, **kw) + +class SendCatchLogTest2(unittest.TestCase): + + def test_error_logged_if_deferred_not_supported(self): + test_signal = object() + test_handler = lambda: defer.Deferred() + log_events = [] + txlog.addObserver(log_events.append) + dispatcher.connect(test_handler, test_signal) + send_catch_log(test_signal) + self.failUnless(log_events) + self.failUnless("Cannot return deferreds from signal handler" in str(log_events)) + txlog.removeObserver(log_events.append) + self.flushLoggedErrors() + dispatcher.disconnect(test_handler, test_signal) diff --git a/scrapy/utils/signal.py b/scrapy/utils/signal.py index 6ea68b0da..d05ab5c70 100644 --- a/scrapy/utils/signal.py +++ b/scrapy/utils/signal.py @@ -1,5 +1,6 @@ """Helper functinos for working with signals""" +from twisted.internet.defer import maybeDeferred, DeferredList, Deferred from twisted.python.failure import Failure from scrapy.xlib.pydispatch.dispatcher import Any, Anonymous, liveReceivers, \ @@ -19,6 +20,9 @@ def send_catch_log(signal=Any, sender=Anonymous, *arguments, **named): try: response = robustApply(receiver, signal=signal, sender=sender, *arguments, **named) + if isinstance(response, Deferred): + log.msg("Cannot return deferreds from signal handler: %s" % \ + receiver, log.ERROR, spider=spider) except dont_log: result = Failure() except Exception: @@ -30,6 +34,30 @@ def send_catch_log(signal=Any, sender=Anonymous, *arguments, **named): responses.append((receiver, result)) return responses +def send_catch_log_deferred(signal=Any, sender=Anonymous, *arguments, **named): + """Like send_catch_log but supports returning deferreds on signal handlers. + Returns a deferred that gets fired once all signal handlers deferreds were + fired. + """ + def logerror(failure, recv): + if dont_log is None or not isinstance(failure.value, dont_log): + log.err(failure, "Error caught on signal handler: %s" % recv, \ + spider=spider) + return failure + + dont_log = named.pop('dont_log', None) + spider = named.get('spider', None) + dfds = [] + for receiver in liveReceivers(getAllReceivers(sender, signal)): + d = maybeDeferred(robustApply, receiver, signal=signal, sender=sender, + *arguments, **named) + d.addErrback(logerror, receiver) + d.addBoth(lambda result: (receiver, result)) + dfds.append(d) + d = DeferredList(dfds) + d.addCallback(lambda out: [x[1] for x in out]) + return d + def disconnect_all(signal=Any, sender=Any): """Disconnect all signal handlers. Useful for cleaning up after running tests