diff --git a/scrapy/command/cmdline.py b/scrapy/command/cmdline.py index c75829371..25f531cea 100644 --- a/scrapy/command/cmdline.py +++ b/scrapy/command/cmdline.py @@ -129,7 +129,7 @@ def execute(argv=None): del args[0] # remove command name from args _save_command_executed(cmdname, cmd, args, opts) from scrapy.core.manager import scrapymanager - scrapymanager.configure() + scrapymanager.configure(control_reactor=True) ret = _run_command(cmd, args, opts) if ret is False: parser.print_help() diff --git a/scrapy/contrib/spidermanager.py b/scrapy/contrib/spidermanager.py index d5ea61284..349e878fe 100644 --- a/scrapy/contrib/spidermanager.py +++ b/scrapy/contrib/spidermanager.py @@ -9,8 +9,6 @@ import urlparse from twisted.plugin import getCache from twisted.python.rebuild import rebuild -from scrapy.xlib.pydispatch import dispatcher -from scrapy.core import signals from scrapy.spider.models import ISpider from scrapy import log from scrapy.conf import settings @@ -25,10 +23,9 @@ class TwistedPluginSpiderManager(object): self.default_domain = None self.force_domain = None self.spider_modules = None - dispatcher.connect(self._domain_closed, signal=signals.domain_closed) - def fromdomain(self, domain_name): - return self.asdict().get(domain_name) + def fromdomain(self, domain): + return self.asdict().get(domain) def fromurl(self, url): if self.force_domain: @@ -112,10 +109,11 @@ class TwistedPluginSpiderManager(object): sys.stderr.write("Interrupted while loading Scrapy spiders\n") sys.exit(2) - def _domain_closed(self, domain, spider): + def close_domain(self, domain): """Reload spider module to release any resources held on to by the spider """ + spider = self._spiders[domain] module_name = spider.__module__ module = sys.modules[module_name] if hasattr(module, 'SPIDER'): diff --git a/scrapy/core/downloader/manager.py b/scrapy/core/downloader/manager.py index 685a206e3..51ff45b08 100644 --- a/scrapy/core/downloader/manager.py +++ b/scrapy/core/downloader/manager.py @@ -5,6 +5,7 @@ Download web pages using asynchronous IO from time import time from twisted.internet import reactor, defer +from twisted.python.failure import Failure from scrapy.core.exceptions import IgnoreRequest from scrapy.spider import spiders @@ -31,8 +32,8 @@ class SiteInfo(object): else: self.max_concurrent_requests = max_concurrent_requests - self.queue = [] self.active = set() + self.queue = [] self.transferring = set() self.closing = False self.lastseen = 0 @@ -56,7 +57,7 @@ class Downloader(object): self.concurrent_domains = settings.getint('CONCURRENT_DOMAINS') def fetch(self, request, spider): - """ Main method to use to request a download + """Main method to use to request a download This method includes middleware mangling. Middleware can returns a Response object, then request never reach downloader queue, and it will @@ -64,7 +65,7 @@ class Downloader(object): """ site = self.sites[spider.domain_name] if site.closing: - raise IgnoreRequest('Can\'t fetch on a closing domain') + raise IgnoreRequest('Cannot fetch on a closing domain') site.active.add(request) def _deactivate(_): @@ -72,25 +73,21 @@ class Downloader(object): self._close_if_idle(spider.domain_name) return _ - return self.middleware.download(self.enqueue, request, spider).addBoth(_deactivate) + dfd = self.middleware.download(self.enqueue, request, spider) + return dfd.addBoth(_deactivate) def enqueue(self, request, spider): """Enqueue a Request for a effective download from site""" - deferred = defer.Deferred() site = self.sites[spider.domain_name] + if site.closing: + raise IgnoreRequest + deferred = defer.Deferred() site.queue.append((request, deferred)) - self.process_queue(spider) + self._process_queue(spider) return deferred - def process_queue(self, spider): - try: - self._process_queue(spider) - except: - log.exc('Downloader process queue bug') - def _process_queue(self, spider): - """ Effective download requests from site queue - """ + """Effective download requests from site queue""" domain = spider.domain_name site = self.sites.get(domain) if not site: @@ -101,14 +98,18 @@ class Downloader(object): if site.download_delay: penalty = site.download_delay - now + site.lastseen if penalty > 0: - reactor.callLater(penalty, self.process_queue, spider=spider) + reactor.callLater(penalty, self._process_queue, spider=spider) return site.lastseen = now - # Process requests in queue if there are free slots to transfer for this site + # Process enqueued requests if there are free slots to transfer for this site while site.queue and site.free_transfer_slots() > 0: request, deferred = site.queue.pop(0) - self._download(site, request, spider).chainDeferred(deferred) + if site.closing: + dfd = defer.fail(Failure(IgnoreRequest())) + else: + dfd = self._download(site, request, spider) + dfd.chainDeferred(deferred) self._close_if_idle(domain) @@ -130,7 +131,13 @@ class Downloader(object): site.transferring.add(request) def finish_transferring(_): site.transferring.remove(request) - self.process_queue(spider) + self._process_queue(spider) + # avoid partially downloaded responses from propagating to the + # downloader middleware, to speed-up the closing process + if site.closing: + log.msg("Crawled while closing domain: %s" % request, \ + level=log.DEBUG) + raise IgnoreRequest return _ return dfd.addBoth(finish_transferring) @@ -153,7 +160,7 @@ class Downloader(object): site.closing = True spider = spiders.fromdomain(domain) - self.process_queue(spider) + self._process_queue(spider) def has_capacity(self): """Does the downloader have capacity to handle more domains""" diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index ba85f3283..f1aa7ec91 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -6,7 +6,7 @@ For more information see docs/topics/architecture.rst """ from datetime import datetime -from twisted.internet import reactor, task +from twisted.internet import reactor, task, defer from twisted.python.failure import Failure from scrapy.xlib.pydispatch import dispatcher @@ -30,8 +30,8 @@ class ExecutionEngine(object): self.keep_alive = False self.closing = {} # dict (domain -> reason) of spiders being closed self.running = False + self.killed = False self.paused = False - self.control_reactor = True self._next_request_pending = set() self._mainloop_task = task.LoopingCall(self._mainloop) @@ -45,34 +45,34 @@ class ExecutionEngine(object): self.scraper = Scraper(self) self.configured = True - def start(self, control_reactor=True): + def start(self): """Start the execution engine""" if self.running: return - self.control_reactor = control_reactor self.start_time = datetime.utcnow() send_catch_log(signal=signals.engine_started, sender=self.__class__) self._mainloop_task.start(5.0, now=True) reactor.callWhenRunning(self._mainloop) self.running = True - if control_reactor: - reactor.run() # blocking call def stop(self): - """Stop the execution engine""" + """Stop the execution engine gracefully""" if not self.running: return self.running = False for domain in self.open_domains: - spider = spiders.fromdomain(domain) - send_catch_log(signal=signals.domain_closed, sender=self.__class__, \ - domain=domain, spider=spider, reason='shutdown') - stats.close_domain(domain, reason='shutdown') + reactor.addSystemEventTrigger('before', 'shutdown', \ + self.close_domain, domain, reason='shutdown') if self._mainloop_task.running: self._mainloop_task.stop() - if self.control_reactor and reactor.running: - reactor.stop() - send_catch_log(signal=signals.engine_stopped, sender=self.__class__) + + def kill(self): + """Forces shutdown without waiting for pending transfers to finish. + stop() must have been called first + """ + if self.running: + return + self.killed = True def pause(self): """Pause the execution engine""" @@ -83,7 +83,8 @@ class ExecutionEngine(object): self.paused = False def is_idle(self): - return self.scheduler.is_idle() and self.downloader.is_idle() and self.scraper.is_idle() + 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() @@ -130,12 +131,15 @@ class ExecutionEngine(object): spider = spiders.fromdomain(domain) dwld = mustbe_deferred(self.download, request, spider) dwld.chainDeferred(deferred).addBoth(lambda _: deferred) - return dwld.addErrback(log.err) + dwld.addErrback(log.err, "Unhandled error on engine._next_request") + return dwld def domain_is_idle(self, domain): - scraper_idle = domain in self.scraper.sites and self.scraper.sites[domain].is_idle() + scraper_idle = domain in self.scraper.sites \ + and self.scraper.sites[domain].is_idle() pending = self.scheduler.domain_has_pending_requests(domain) - downloading = domain in self.downloader.sites and self.downloader.sites[domain].active + downloading = domain in self.downloader.sites \ + and self.downloader.sites[domain].active return scraper_idle and not (pending or downloading) def domain_is_closed(self, domain): @@ -155,7 +159,7 @@ class ExecutionEngine(object): def crawl(self, request, spider): schd = mustbe_deferred(self.schedule, request, spider) schd.addBoth(self.scraper.enqueue_scrape, request, spider) - schd.addErrback(log.err) + schd.addErrback(log.err, "Unhandled error on engine.crawl()") schd.addBoth(lambda _: self.next_request(spider.domain_name)) def schedule(self, request, spider): @@ -190,8 +194,8 @@ class ExecutionEngine(object): assert isinstance(response, (Response, Request)) if isinstance(response, Response): response.request = request # tie request to response received - log.msg("Crawled %s (referer: <%s>)" % (response, referer), level=log.DEBUG, \ - domain=domain) + log.msg("Crawled %s (referer: <%s>)" % (response, referer), \ + level=log.DEBUG, domain=domain) return response elif isinstance(response, Request): newrequest = response @@ -201,12 +205,17 @@ class ExecutionEngine(object): def _on_error(_failure): """handle an error processing a page""" - ex = _failure.value - errmsg = str(_failure) if not isinstance(ex, IgnoreRequest) \ - else _failure.getErrorMessage() - log.msg("Downloading <%s> (referer: <%s>): %s" % (request.url, referer, errmsg), \ - log.ERROR, domain=domain) - return Failure(IgnoreRequest(str(ex))) + exc = _failure.value + if isinstance(exc, IgnoreRequest): + errmsg = _failure.getErrorMessage() + level = exc.level + else: + errmsg = str(_failure) + level = log.ERROR + if errmsg: + log.msg("Downloading <%s> (referer: <%s>): %s" % (request.url, \ + referer, errmsg), level=level, domain=domain) + return Failure(IgnoreRequest(str(exc))) def _on_complete(_): self.next_request(domain) @@ -249,7 +258,7 @@ class ExecutionEngine(object): self.next_request(domain) return except: - log.err(_why="Exception catched on domain_idle signal dispatch") + log.err("Exception catched on domain_idle signal dispatch") if self.domain_is_idle(domain): self.close_domain(domain, reason='finished') @@ -265,14 +274,19 @@ class ExecutionEngine(object): self.closing[domain] = reason self.downloader.close_domain(domain) self.scheduler.clear_pending_requests(domain) - self._finish_closing_domain_if_idle(domain) + return self._finish_closing_domain_if_idle(domain) + return defer.succeed(None) def _finish_closing_domain_if_idle(self, domain): """Call _finish_closing_domain if domain is idle""" - if self.domain_is_idle(domain): + if self.domain_is_idle(domain) or self.killed: self._finish_closing_domain(domain) else: - reactor.callLater(5, self._finish_closing_domain_if_idle, domain) + dfd = defer.Deferred() + dfd.addCallback(self._finish_closing_domain_if_idle) + delay = 5 if self.running else 1 + reactor.callLater(delay, dfd.callback, domain) + return dfd def _finish_closing_domain(self, domain): """This function is called after the domain has been closed""" @@ -284,6 +298,10 @@ class ExecutionEngine(object): domain=domain, spider=spider, reason=reason) stats.close_domain(domain, reason=reason) log.msg("Domain closed (%s)" % reason, domain=domain) - self._mainloop() + spiders.close_domain(domain) + if self.running: + self._mainloop() + elif not self.open_domains: + send_catch_log(signal=signals.engine_stopped, sender=self.__class__) scrapyengine = ExecutionEngine() diff --git a/scrapy/core/exceptions.py b/scrapy/core/exceptions.py index e1fbd6f85..97e39de14 100644 --- a/scrapy/core/exceptions.py +++ b/scrapy/core/exceptions.py @@ -5,6 +5,8 @@ These exceptions are documented in docs/topics/exceptions.rst. Please don't add new exceptions here without documenting them there. """ +from scrapy import log + # Internal class NotConfigured(Exception): @@ -15,7 +17,13 @@ class NotConfigured(Exception): class IgnoreRequest(Exception): """Indicates a decision was made not to process a request""" - pass + + def __init__(self, msg='', level=log.ERROR): + self.msg = msg + self.level = level + + def __str__(self): + return self.msg class DontCloseDomain(Exception): """Request the domain not to be closed yet""" diff --git a/scrapy/core/manager.py b/scrapy/core/manager.py index be7de69da..dcef06e53 100644 --- a/scrapy/core/manager.py +++ b/scrapy/core/manager.py @@ -9,7 +9,48 @@ from scrapy.core.engine import scrapyengine from scrapy.spider import spiders from scrapy.utils.misc import arg_to_iter from scrapy.utils.url import is_url -from scrapy.conf import settings +from scrapy.utils.ossignal import install_shutdown_handlers, signal_names + +def _parse_args(args): + """Parse crawl arguments and return a dict of domains -> list of requests""" + requests, urls, sites = set(), set(), set() + for a in args: + if isinstance(a, Request): + requests.add(a) + elif is_url(a): + urls.add(a) + else: + sites.add(a) + + perdomain = {} + + # sites + for domain in sites: + spider = spiders.fromdomain(domain) + if not spider: + log.msg('Could not find spider for %s' % domain, log.ERROR) + continue + reqs = spider.start_requests() + perdomain.setdefault(domain, []).extend(reqs) + + # urls + for url in urls: + spider = spiders.fromurl(url) + if spider: + for req in arg_to_iter(spider.make_requests_from_url(url)): + perdomain.setdefault(spider.domain_name, []).append(req) + else: + log.msg('Could not find spider for <%s>' % url, log.ERROR) + + # requests + for request in requests: + spider = spiders.fromurl(request.url) + if not spider: + log.msg('Could not find spider for %s' % request, log.ERROR) + continue + perdomain.setdefault(spider.domain_name, []).append(request) + return perdomain + class ExecutionManager(object): """Process a list of sites or urls. @@ -27,8 +68,8 @@ class ExecutionManager(object): def configure(self, control_reactor=True): self.control_reactor = control_reactor if control_reactor: - self._install_signals() - reactor.addSystemEventTrigger('before', 'shutdown', self.stop) + install_shutdown_handlers(self._signal_shutdown) + reactor.addSystemEventTrigger('before', 'shutdown', scrapyengine.stop) if not log.started: log.start() @@ -45,7 +86,7 @@ class ExecutionManager(object): def crawl(self, *args): """Schedule the given args for crawling. args is a list of urls or domains""" - requests = self._parse_args(args) + requests = _parse_args(args) # schedule initial requests to be scraped at engine start for domain in requests or (): spider = spiders.fromdomain(domain) @@ -54,88 +95,40 @@ class ExecutionManager(object): def runonce(self, *args): """Run the engine until it finishes scraping all domains and then exit""" - if not self.configured: - self.configure() + assert self.configured, "Scrapy Manger not yet configured" self.crawl(*args) scrapyengine.start() + if self.control_reactor: + reactor.run(installSignalHandlers=False) - def start(self, control_reactor=True): + def start(self): """Start the scrapy server, without scheduling any domains""" - self.configure(control_reactor) + assert self.configured, "Scrapy Manger not yet configured" scrapyengine.keep_alive = True - scrapyengine.start(control_reactor=control_reactor) - if control_reactor: - self.stop() + scrapyengine.start() + if self.control_reactor: + reactor.run(installSignalHandlers=False) def stop(self): """Stop the scrapy server, shutting down the execution engine""" self.interrupted = True scrapyengine.stop() - log.log_level = -999 # disable logging - if self.control_reactor: - signal.signal(signal.SIGTERM, signal.SIG_IGN) - signal.signal(signal.SIGINT, signal.SIG_IGN) - if hasattr(signal, "SIGBREAK"): - signal.signal(signal.SIGBREAK, signal.SIG_IGN) + if self.control_reactor and reactor.running: + reactor.stop() - def reload_spiders(self): - """Reload all enabled spiders except for the ones that are currently - running. - """ - spiders.reload(skip_domains=scrapyengine.open_domains) + def _signal_shutdown(self, signum, _): + signame = signal_names[signum] + log.msg("Received %s, shutting down gracefully. Send again to force " \ + "unclean shutdown" % signame, level=log.INFO) + reactor.callFromThread(self.stop) + install_shutdown_handlers(self._signal_kill) - def _install_signals(self): - def sig_handler_terminate(signalinfo, param): - log.msg('Received shutdown request, waiting for deferreds to finish...', log.INFO) - self.stop() - - signal.signal(signal.SIGTERM, sig_handler_terminate) - # only handle SIGINT if there isn't already a handler (e.g. for Pdb) - if signal.getsignal(signal.SIGINT) == signal.default_int_handler: - signal.signal(signal.SIGINT, sig_handler_terminate) - # Catch Ctrl-Break in windows - if hasattr(signal, "SIGBREAK"): - signal.signal(signal.SIGBREAK, sig_handler_terminate) - - def _parse_args(self, args): - """Parse crawl arguments and return a dict of domains -> list of - requests""" - requests, urls, sites = set(), set(), set() - for a in args: - if isinstance(a, Request): - requests.add(a) - elif is_url(a): - urls.add(a) - else: - sites.add(a) - - perdomain = {} - - # sites - for domain in sites: - spider = spiders.fromdomain(domain) - if not spider: - log.msg('Could not find spider for %s' % domain, log.ERROR) - continue - reqs = spider.start_requests() - perdomain.setdefault(domain, []).extend(reqs) - - # urls - for url in urls: - spider = spiders.fromurl(url) - if spider: - for req in arg_to_iter(spider.make_requests_from_url(url)): - perdomain.setdefault(spider.domain_name, []).append(req) - else: - log.msg('Could not find spider for <%s>' % url, log.ERROR) - - # requests - for request in requests: - spider = spiders.fromurl(request.url) - if not spider: - log.msg('Could not find spider for %s' % request, log.ERROR) - continue - perdomain.setdefault(spider.domain_name, []).append(request) - return perdomain + def _signal_kill(self, signum, _): + signame = signal_names[signum] + log.msg('Received %s twice, forcing unclean shutdown' % signame, \ + level=log.INFO) + log.log_level = log.SILENT # disable logging of confusing tracebacks + reactor.callFromThread(scrapyengine.kill) + install_shutdown_handlers(signal.SIG_IGN) scrapymanager = ExecutionManager() diff --git a/scrapy/log.py b/scrapy/log.py index 4e75854fa..ecb7ef7d4 100644 --- a/scrapy/log.py +++ b/scrapy/log.py @@ -56,19 +56,24 @@ def start(logfile=None, loglevel=None, logstdout=None): def msg(message, level=INFO, component=BOT_NAME, domain=None): """Log message according to the level""" + if level > log_level: + return dispatcher.send(signal=logmessage_received, message=message, level=level, \ domain=domain) system = domain if domain else component - if level <= log_level: - msg_txt = unicode_to_str("%s: %s" % (level_names[level], message)) - log.msg(msg_txt, system=system) + msg_txt = unicode_to_str("%s: %s" % (level_names[level], message)) + log.msg(msg_txt, system=system) def exc(message, level=ERROR, component=BOT_NAME, domain=None): message = message + '\n' + format_exc() msg(message, level, component, domain) -def err(*args, **kwargs): +def err(_stuff=None, _why=None, **kwargs): + if ERROR > log_level: + return domain = kwargs.pop('domain', None) component = kwargs.pop('component', BOT_NAME) kwargs['system'] = domain if domain else component - log.err(*args, **kwargs) + if _why: + _why = unicode_to_str("ERROR: %s" % _why) + log.err(_stuff, _why, **kwargs) diff --git a/scrapy/utils/ossignal.py b/scrapy/utils/ossignal.py new file mode 100644 index 000000000..df4eee5ec --- /dev/null +++ b/scrapy/utils/ossignal.py @@ -0,0 +1,28 @@ + +from __future__ import absolute_import + +from twisted.internet import reactor + +import signal + +signal_names = {} +for signame in dir(signal): + if signame.startswith("SIG"): + signum = getattr(signal, signame) + if isinstance(signum, int): + signal_names[signum] = signame + +def install_shutdown_handlers(function, override_sigint=True): + """Install the given function as a signal handler for all common shutdown + signals (such as SIGINT, SIGTERM, etc). If override_sigint is ``False`` the + SIGINT handler won't be install if there is already a handler in place + (e.g. Pdb) + """ + reactor._handleSignals() + signal.signal(signal.SIGTERM, function) + if signal.getsignal(signal.SIGINT) == signal.default_int_handler or \ + override_sigint: + signal.signal(signal.SIGINT, function) + # Catch Ctrl-Break in windows + if hasattr(signal, "SIGBREAK"): + signal.signal(signal.SIGBREAK, function)