From c913d6ffcf922f3dde5592b35c95a1813340a27b Mon Sep 17 00:00:00 2001 From: Pablo Hoffman Date: Thu, 23 Oct 2008 00:37:08 +0000 Subject: [PATCH] added generic pre-run hooks to cluster workers to deocuple pysvn from worker code --HG-- extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%40329 --- .../contrib/pbcluster/hooks/__init__.py | 13 +++++++ .../scrapy/contrib/pbcluster/hooks/svn.py | 23 +++++++++++ .../contrib/pbcluster/worker/manager.py | 39 ++++++++++--------- 3 files changed, 57 insertions(+), 18 deletions(-) create mode 100644 scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/hooks/__init__.py create mode 100644 scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/hooks/svn.py diff --git a/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/hooks/__init__.py b/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/hooks/__init__.py new file mode 100644 index 000000000..a66f00650 --- /dev/null +++ b/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/hooks/__init__.py @@ -0,0 +1,13 @@ +""" +This module contains pre-run hooks that can be attached to scrapy workers. + +Pre-run hooks must be callable objects (ie. functions) which implement this +interface: + +pre_hook(domain, spider_settings) + +domain is the domain to be scraped +spider_settings is the settings to use to scrape it + +Values returned from the pre-run hooks will be ignored. +""" diff --git a/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/hooks/svn.py b/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/hooks/svn.py new file mode 100644 index 000000000..b3f7a6d6c --- /dev/null +++ b/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/hooks/svn.py @@ -0,0 +1,23 @@ +""" +This module contains hooks for updating code via svn. Useful for running before +starting to crawl a domain, for example to update spider code +""" + +import pysvn + +from scrapy.conf import settings +from scrapy import log + +SVN_DIR = settings['SVN_DIR'] +SVN_USER = settings['SVN_USER'] +SVN_PASS = settings['SVN_PASS'] + +def svnup(domain, spider_settings): + c = pysvn.Client() + c.callback_get_login = lambda x,y,z: (True, SVN_USER, SVN_PASS, False) + try: + r = c.update(SVN_DIR) + log.msg("ClusterWorker: SVN code updated to revision %s (triggered by domain %s)" % \ + (r[0].number, domain), level=log.DEBUG) + except pysvn.ClientError, e: + log.msg("ClusterWorker: unable to update svn code - %s" % e, level=log.WARNING) diff --git a/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/worker/manager.py b/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/worker/manager.py index 6bff3bdc8..eb01bef02 100644 --- a/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/worker/manager.py +++ b/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/worker/manager.py @@ -5,12 +5,14 @@ import datetime import cPickle as pickle from twisted.internet import protocol, reactor +from twisted.internet.error import ProcessDone from twisted.spread import pb from scrapy import log -from scrapy.core.exceptions import NotConfigured -from scrapy.conf import settings from scrapy.core.engine import scrapyengine +from scrapy.core.exceptions import NotConfigured +from scrapy.utils.misc import load_class +from scrapy.conf import settings class ScrapyProcessProtocol(protocol.ProcessProtocol): @@ -70,9 +72,14 @@ class ScrapyProcessProtocol(protocol.ProcessProtocol): self.status = "running" self.procman.update_master(self.domain, "running") - def processEnded(self, reason): - log.msg("ClusterWorker: finished domain=%s, pid=%d, log=%s" % (self.domain, self.pid, self.logfile)) - log.msg("Reason type: %s. value: %s" % (reason.type, reason.value) ) + def processEnded(self, status): + if isinstance(status.value, ProcessDone): + st = "done" + er = "" + else: + st = "terminated" + er = ", error=%s" % str(status.value) + log.msg("ClusterWorker: finished domain=%s, status=%s, pid=%d, log=%s%s" % (self.domain, st, self.pid, self.logfile, er)) del self.procman.running[self.domain] del self.procman.crawlers[self.pid] self.procman.update_master(self.domain, "scraped") @@ -88,6 +95,7 @@ class ClusterWorker(pb.Root): self.running = {} # dict of domain->ScrapyProcessControl self.crawlers = {} # dict of pid->scrapy process remote pb connection self.starttime = datetime.datetime.utcnow() + self.prerun_hooks = [load_class(f) for f in settings.getlist('CLUSTER_WORKER_PRERUN_HOOKS', [])] port = settings.getint('CLUSTER_WORKER_PORT') scrapyengine.listenTCP(port, pb.PBServerFactory(self)) log.msg("Using sys.path: %s" % repr(sys.path), level=log.DEBUG) @@ -150,7 +158,7 @@ class ClusterWorker(pb.Root): d.addCallbacks(callback=_close, errback=lambda reason: log.msg(reason, log.ERROR)) return self.status(0, "Stopped process %s" % proc) else: - return self.status(1, "%s: domain not running." % domain) + return self.status(1, "%s: domain not running" % domain) def remote_status(self): """Return worker status as a dict. For infomation about the keys see @@ -167,21 +175,16 @@ class ClusterWorker(pb.Root): scrapy_proc = ScrapyProcessProtocol(self, domain, logfile, spider_settings) args = [sys.executable, sys.argv[0], 'crawl', domain] self.running[domain] = scrapy_proc - try: - import pysvn - c = pysvn.Client() - r = c.update(settings.get("CLUSTER_WORKER_SVNWORKDIR", ".")) - log.msg("Updated to revision %s." %r[0].number, level=log.DEBUG) - except pysvn.ClientError, e: - log.msg("Unable to svn update: %s" % e, level=log.WARNING) - except ImportError: - log.msg("pysvn module not available.", level=log.WARNING) + + for prerun_hook in self.prerun_hooks: + prerun_hook(domain, spider_settings) + reactor.spawnProcess(scrapy_proc, sys.executable, args=args, env=scrapy_proc.env) - return self.status(0, "Started process %s." % scrapy_proc) + return self.status(0, "Started process %s" % scrapy_proc) else: - return self.status(2, "Domain %s already running." % domain ) + return self.status(2, "Domain %s already running" % domain ) else: - return self.status(1, "No free slot to run another process.") + return self.status(1, "No free slot to run another process") def remote_register_crawler(self, pid, crawler): """Register the crawler to the list of crawlers managed by this worker"""