diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index e1050f084..fbc3ec010 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -5,7 +5,6 @@ For more information see docs/topics/architecture.rst """ from datetime import datetime -from itertools import imap from twisted.internet import reactor, task from twisted.internet.error import CannotListenError @@ -13,15 +12,14 @@ from twisted.python.failure import Failure from pydispatch import dispatcher from scrapy import log -from scrapy.stats import stats from scrapy.conf import settings from scrapy.core import signals from scrapy.core.scheduler import Scheduler from scrapy.core.downloader import Downloader +from scrapy.core.scraper import Scraper from scrapy.core.exceptions import IgnoreRequest, DontCloseDomain from scrapy.http import Response, Request from scrapy.spider import spiders -from scrapy.spider.middleware import SpiderMiddlewareManager from scrapy.utils.misc import load_object from scrapy.utils.defer import mustbe_deferred @@ -45,8 +43,7 @@ class ExecutionEngine(object): self.scheduler = scheduler or Scheduler() self.domain_scheduler = load_object(settings['DOMAIN_SCHEDULER'])() self.downloader = downloader or Downloader() - self.spidermiddleware = SpiderMiddlewareManager() - self._scraping = {} + self.scraper = Scraper(self) self.configured = True def addtask(self, function, interval, args=None, kwargs=None, now=False): @@ -136,7 +133,7 @@ class ExecutionEngine(object): self.paused = False def is_idle(self): - return self.scheduler.is_idle() and self.downloader.is_idle() and not self._scraping + return self.scheduler.is_idle() and self.downloader.is_idle() and self.scraper.is_idle() def next_domain(self): domain = self.domain_scheduler.next_domain() @@ -179,7 +176,7 @@ class ExecutionEngine(object): self._domain_idle(domain) def domain_is_idle(self, domain): - scraping = self._scraping.get(domain) + scraping = domain in self.scraper.sites and bool(self.scraper.sites[domain]) pending = self.scheduler.domain_has_pending_requests(domain) downloading = domain in self.downloader.sites and self.downloader.sites[domain].active return not (pending or downloading or scraping) @@ -199,56 +196,10 @@ class ExecutionEngine(object): return self.downloader.sites.keys() def crawl(self, request, spider): - domain = spider.domain_name - - def _process_response(response): - assert isinstance(response, (Response, Exception)), \ - "Expecting Response or Exception, got %s" % type(response).__name__ - - def cb_spidermiddleware_output(spmw_result): - def cb_spider_output(output): - if domain in self.closing: - return - elif isinstance(output, Request): - signals.send_catch_log(signal=signals.request_received, sender=self.__class__, request=output, spider=spider, response=response) - self.crawl(request=output, spider=spider) - 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) - - return task.coiterate(imap(cb_spider_output, spmw_result)) - - def eb_user(_failure): - if not isinstance(_failure.value, IgnoreRequest): - referer = None if not isinstance(response, Response) else response.request.headers.get('Referer', None) - log.msg("Error while spider was processing <%s> from <%s>: %s" % (request.url, referer, _failure), log.ERROR, domain=domain) - stats.incpath("%s/spider_exceptions/%s" % (domain, _failure.value.__class__.__name__)) - - scd = mustbe_deferred(self.spidermiddleware.scrape, request, response, spider) - scd.addCallbacks(cb_spidermiddleware_output, eb_user) - - self._scraping[domain].add(response) - def _remove(_): - self._scraping[domain].remove(response) - self.next_request(domain) - return _ - - scd.addBoth(_remove) - scd.addErrback(log.err, 'FRAMEWORK BUG processing %s' % request, domain=domain) - return scd - - def _cleanfailure(_failure): - ex = _failure.value - if not isinstance(ex, IgnoreRequest): - log.msg("Unknown error propagated in %s: %s" % (request, _failure), log.ERROR, domain=domain) - request.deferred.addErrback(lambda _:None) - request.deferred.errback(_failure) # TODO: merge into spider middleware. - schd = mustbe_deferred(self.schedule, request, spider) - schd.addCallbacks(_process_response, _cleanfailure) - return schd.addErrback(log.err) + schd.addBoth(self.scraper.scrape, request, spider) + schd.addErrback(log.err) + schd.addBoth(lambda _: self.next_request(spider.domain_name)) def schedule(self, request, spider): domain = spider.domain_name @@ -313,7 +264,7 @@ class ExecutionEngine(object): self.next_request(domain) self.downloader.open_domain(domain) - self._scraping[domain] = set() + self.scraper.open_domain(domain) signals.send_catch_log(signals.domain_open, sender=self.__class__, domain=domain, spider=spider) signals.send_catch_log(signals.domain_opened, sender=self.__class__, domain=domain, spider=spider) @@ -362,7 +313,7 @@ class ExecutionEngine(object): """This function is called after the domain has been closed""" spider = spiders.fromdomain(domain) self.scheduler.close_domain(domain) - del self._scraping[domain] + self.scraper.close_domain(domain) reason = self.closing.pop(domain, 'finished') signals.send_catch_log(signal=signals.domain_closed, sender=self.__class__, domain=domain, spider=spider, reason=reason) log.msg("Domain closed (%s)" % reason, domain=domain) @@ -382,7 +333,8 @@ class ExecutionEngine(object): "self.downloader.is_idle()", "len(self.downloader.sites)", "self.downloader.has_capacity()", - "len(self._scraping)", + "self.scraper.is_idle()", + "len(self.scraper.sites)", ] domain_tests = [ "self.domain_is_idle(domain)", @@ -394,7 +346,7 @@ class ExecutionEngine(object): "len(self.downloader.sites[domain].transferring)", "self.downloader.sites[domain].closing", "self.downloader.sites[domain].lastseen", - "len(self._scraping[domain])", + "len(self.scraper.sites[domain])", ] for test in global_tests: diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py new file mode 100644 index 000000000..ad239493a --- /dev/null +++ b/scrapy/core/scraper.py @@ -0,0 +1,107 @@ +"""Extract information from pages""" + +from itertools import imap +from twisted.internet import task +from twisted.python.failure import Failure + +from scrapy.utils.defer import defer_result +from scrapy.utils.misc import arg_to_iter +from scrapy.core.exceptions import IgnoreRequest +from scrapy.core import signals +from scrapy.http import Request, Response +from scrapy.spider.middleware import SpiderMiddlewareManager +from scrapy import log +from scrapy.stats import stats + +class Scraper(object): + + def __init__(self, engine): + self.sites = {} + self.middleware = SpiderMiddlewareManager() + self.engine = engine + + def open_domain(self, domain): + """Open the given domain for scraping and allocate resources for it""" + if domain in self.sites: + raise RuntimeError('Scraper domain already opened: %s' % domain) + self.sites[domain] = set() + + def close_domain(self, domain): + """Close a domain being scraped and release its resources""" + if domain not in self.sites: + raise RuntimeError('Scraper domain already closed: %s' % domain) + del self.sites[domain] + + def is_idle(self): + """Return True if there isn't any more spiders to process""" + return not self.sites + + def scrape(self, response, request, spider): + """Handle the downloaded response or failure trough the spider + callback/errback""" + assert isinstance(response, (Response, Failure)) + domain = spider.domain_name + + self.sites[domain].add(response) + def _finish_scraping(_): + self.sites[domain].remove(response) + return _ + + dfd = self._scrape(response, request, spider) # returns spiders processed output + dfd.addErrback(self.handle_spider_error, request, spider) + dfd.addCallback(self.handle_spider_output, request, spider) + dfd.addBoth(_finish_scraping) + dfd.addErrback(log.err, 'Scraper bug processing %s' % request, \ + domain=domain) + return dfd + + def _scrape(self, request_result, request, spider): + """Handle the diferent cases of request's result been a Response or a Failure""" + if not isinstance(request_result, Failure): + return self.middleware.scrape_response(self.call_spider, \ + request_result, request, spider) + else: + # FIXME: don't ignore errors in spider middleware + dfd = self.call_spider(request_result, request, spider) + return dfd.addErrback(self._check_propagated_failure, \ + request_result, request, spider) + + def call_spider(self, result, request, spider): + defer_result(result).chainDeferred(request.deferred) + return request.deferred.addCallback(arg_to_iter) + + def handle_spider_error(self, _failure, request, spider, propagated_failure=None): + referer = request.headers.get('Referer', None) + msg = "SPIDER BUG processing <%s> from <%s>: %s" % (request.url, referer, _failure) + log.msg(msg, log.ERROR, domain=spider.domain_name) + stats.incpath("%s/spider_exceptions/%s" % (spider.domain_name, _failure.value.__class__.__name__)) + + def handle_spider_output(self, result, request, spider): + func = lambda o: self.process_spider_output(o, request, spider) + return task.coiterate(imap(func, result or [])) + + def process_spider_output(self, output, request, spider): + if spider.domain_name in self.engine.closing: + return + elif isinstance(output, Request): + signals.send_catch_log(signal=signals.request_received, request=output, spider=spider) + self.engine.crawl(request=output, spider=spider) + 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=spider.domain_name) + + def _check_propagated_failure(self, spider_failure, propagated_failure, request, spider): + """Log and silence the bugs raised outside of spiders, but still allow + spiders to be notified about general failures while downloading spider + generated requests + """ + # ignored requests are commonly propagated exceptions safes to be silenced + if isinstance(spider_failure.value, IgnoreRequest): + return + elif spider_failure is propagated_failure: + log.err(spider_failure, 'Unhandled error propagated to spider and wasn\'t handled') + return # stop propagating this error + else: + return spider_failure # exceptions raised in the spider code diff --git a/scrapy/spider/middleware.py b/scrapy/spider/middleware.py index 4df7599f3..ffb1d0770 100644 --- a/scrapy/spider/middleware.py +++ b/scrapy/spider/middleware.py @@ -8,9 +8,9 @@ docs/topics/spider-middleware.rst from scrapy import log from scrapy.core.exceptions import NotConfigured -from scrapy.utils.misc import load_object, arg_to_iter -from scrapy.utils.defer import mustbe_deferred, defer_result +from scrapy.utils.misc import load_object from scrapy.utils.middleware import build_middleware_list +from scrapy.utils.defer import mustbe_deferred from scrapy.http import Request from scrapy.conf import settings @@ -55,7 +55,7 @@ class SpiderMiddlewareManager(object): level=log.DEBUG) self.loaded = True - def scrape(self, request, response, spider): + def scrape_response(self, scrape_func, response, request, spider): fname = lambda f:'%s.%s' % (f.im_self.__class__.__name__, f.im_func.__name__) def process_spider_input(response): @@ -66,7 +66,7 @@ class SpiderMiddlewareManager(object): (fname(method), type(result)) if result is not None: return result - return self.call(request, response, spider) + return scrape_func(response, request, spider) def process_spider_exception(_failure): exception = _failure.value @@ -93,11 +93,6 @@ class SpiderMiddlewareManager(object): dfd.addCallback(process_spider_output) return dfd - def call(self, request, result, spider): - defer_result(result).chainDeferred(request.deferred) - request.deferred.addCallback(arg_to_iter) - return request.deferred - def _validate_output(self, request, result, spider): """Every request returned by spiders must be instanciate with a callback""" for r in result: