diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py new file mode 100644 index 000000000..8d66b9799 --- /dev/null +++ b/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py @@ -0,0 +1,33 @@ +from twisted.spread import pb +from twisted.internet import reactor + +from scrapy.conf import settings +from scrapy.core import log +from scrapy.core.manager import scrapymanager + +class Broker(pb.Referenceable): + def __init__(self, crawler, remote): + self.__remote = remote + self.__crawler = crawler + try: + deferred = self.__remote.callRemote("register_crawler", domain, self) + except pb.DeadReferenceError: + self._set_status(None) + log.msg("Lost connection to node %s." % (self.name), log.ERROR) + else: + deferred.addCallbacks(callback=self._set_status, errback=lambda reason: log.msg(reason, log.ERROR)) + def remote_stop(self): + scrapymanager.stop() + +class ClusterCrawler: + def __init__(self): + if not settings.getbool('CLUSTER_CRAWLER_ENABLED'): + raise NotConfigured + + self.worker = None + + factory = pb.PBClientFactory() + reactor.connectTCP("localhost", settings.getint('CLUSTER_WORKER_PORT'), factory) + d = factory.getRootObject() + d.addCallbacks(callback=lambda obj: self.worker=Node(self, obj), errback=lambda reason: log.msg(reason, log.ERROR)) + \ No newline at end of file diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py index 20737275c..a5aa97fe0 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py @@ -27,7 +27,7 @@ def my_import(name): mod = getattr(mod, comp) return mod -class Node(pb.Referenceable): +class Broker(pb.Referenceable): def __init__(self, remote, name, master): self.__remote = remote self.alive = False @@ -177,7 +177,7 @@ class ClusterMaster: self.loading = [] self.nodes = {} self.start_time = datetime.datetime.utcnow() - #on how statistics works, see self.update_nodes() and Nodes.remote_update() + #on how statistics works, see self.update_nodes() and Broker.remote_update() self.statistics = {"domains": {"running": set(), "scraped": {}, "lost_count": {}, "lost": set()}, "scraped_count": 0 } self.global_settings = {} #load cluster global settings @@ -231,7 +231,7 @@ class ClusterMaster: def add_node(self, cworker, name): """Add node given its node""" - node = Node(cworker, name, self) + node = Broker(cworker, name, self) self.nodes[name] = node log.msg("Added cluster worker %s" % name) diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py index 981e1263e..c3d2f7502 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py @@ -19,7 +19,7 @@ class ScrapyProcessProtocol(protocol.ProcessProtocol): self.env = {} #We conserve original setting format for info purposes (avoid lots of unnecesary "SCRAPY_") self.scrapy_settings = spider_settings or {} - self.scrapy_settings.update({'LOGFILE': self.logfile, 'CLUSTER_WORKER_ENABLED': '0', 'WEBCONSOLE_ENABLED': '0'}) + self.scrapy_settings.update({'LOGFILE': self.logfile, 'CLUSTER_WORKER_ENABLED': 0, 'CLUSTER_CRAWLER_ENABLED': 1, 'WEBCONSOLE_ENABLED': 0}) pickled_settings = pickle.dumps(self.scrapy_settings) self.env["SCRAPY_PICKLED_SETTINGS_TO_OVERRIDE"] = pickled_settings self.env["PYTHONPATH"] = ":".join(sys.path)#this is need so this crawl process knows where to locate local_scrapy_settings. @@ -41,6 +41,7 @@ class ScrapyProcessProtocol(protocol.ProcessProtocol): 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) ) del self.procman.running[self.domain] + del self.procman.crawlers[self.domain] self.procman.update_master(self.domain, "scraped") class ClusterWorker(pb.Root): @@ -51,7 +52,8 @@ class ClusterWorker(pb.Root): self.maxproc = settings.getint('CLUSTER_WORKER_MAXPROC') self.logdir = settings['CLUSTER_LOGDIR'] - self.running = {} + self.running = {}#a dict domain->ScrapyProcessControl + self.crawlers = {}#a dict domain->scrapy process remote pb connection self.starttime = datetime.datetime.utcnow() port = settings.getint('CLUSTER_WORKER_PORT') scrapyengine.listenTCP(port, pb.PBServerFactory(self)) @@ -86,8 +88,8 @@ class ClusterWorker(pb.Root): if domain in self.running: proc = self.running[domain] log.msg("ClusterWorker: Sending shutdown signal to domain=%s, pid=%d" % (domain, proc.pid)) - proc.transport.signalProcess('INT') - proc.status = "closing" + d = self.crawler["domain"].callRemote("stop") + d.addCallbacks(callback=lambda x: proc.status="closing", 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) @@ -118,3 +120,6 @@ class ClusterWorker(pb.Root): return self.status(0, "Started process %s." % scrapy_proc) return self.status(2, "Domain %s already running." % domain ) return self.status(1, "No free slot to run another process.") + + def remote_register_crawler(self, domain, crawler): + self.crawlers['domain'] = crawler \ No newline at end of file