diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 8ea50a73d..100f519ef 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -23,7 +23,6 @@ from scrapy.item import ScrapedItem from scrapy.item.pipeline import ItemPipelineManager from scrapy.spider import spiders from scrapy.spider.middleware import SpiderMiddlewareManager -from scrapy.utils.defer import chain_deferred, deferred_imap from scrapy.utils.misc import load_object from scrapy.utils.defer import mustbe_deferred @@ -219,7 +218,6 @@ class ExecutionEngine(object): signals.send_catch_log(signal=signals.item_dropped, sender=self.__class__, item=item, spider=spider, response=response, exception=pipe_result.value) else: signals.send_catch_log(signal=signals.item_passed, sender=self.__class__, item=item, spider=spider, response=response, pipe_output=pipe_result) - self.next_request(spider) if domain in self.closing: return @@ -234,9 +232,10 @@ class ExecutionEngine(object): elif output is None: pass # may be next time. else: - log.msg("Spider must return Request, ScrapedItem or None, got '%s' while processing %s" % (type(output).__name__, request), log.WARNING, domain=domain) + log.msg("Spider must return Request, ScrapedItem or None, got '%s' while processing %s" \ + % (type(output).__name__, request), log.WARNING, domain=domain) - return deferred_imap(cb_spider_output, spmw_result) + return task.coiterate(cb_spider_output, spmw_result) def eb_user(_failure): if not isinstance(_failure.value, IgnoreRequest): @@ -248,7 +247,7 @@ class ExecutionEngine(object): log.msg('FRAMEWORK BUG processing %s: %s' % (request, _failure), log.ERROR, domain=domain) - scd = self.spidermiddleware.scrape(request, response, spider) + scd = mustbe_deferred(self.spidermiddleware.scrape, request, response, spider) scd.addCallbacks(cb_spidermiddleware_output, eb_user) scd.addErrback(eb_framework) @@ -263,7 +262,7 @@ class ExecutionEngine(object): request.deferred.addErrback(lambda _:None) request.deferred.errback(_failure) # TODO: merge into spider middleware. - schd = self.schedule(request, spider) + schd = mustbe_deferred(self.schedule, request, spider) schd.addCallbacks(_process_response, _cleanfailure) return schd @@ -304,9 +303,9 @@ class ExecutionEngine(object): return response elif isinstance(response, Request): newrequest = response - schd = self.schedule(newrequest, spider) - chain_deferred(schd, newrequest.deferred) - return schd + schd = mustbe_deferred(self.schedule, newrequest, spider) + schd.chainDeferred(newrequest.deferred) + return newrequest.deferred def _on_error(_failure): """handle an error processing a page""" @@ -317,13 +316,12 @@ class ExecutionEngine(object): def _on_complete(_): self.next_request(spider) + return _ - dwld = self.downloader.fetch(request, spider) + dwld = mustbe_deferred(self.downloader.fetch, request, spider) dwld.addCallbacks(_on_success, _on_error) - deferred = defer.Deferred() - chain_deferred(dwld, deferred) dwld.addBoth(_on_complete) - return deferred + return dwld def open_domain(self, domain): log.msg("Domain opened", domain=domain)