diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/__init__.py b/scrapy/trunk/scrapy/contrib/pbcluster/__init__.py index 40aa182de..84b6df381 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/__init__.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/__init__.py @@ -1,2 +1,3 @@ from scrapy.contrib.pbcluster.worker.manager import ClusterWorker -from scrapy.contrib.pbcluster.master.web import ClusterMasterWeb \ No newline at end of file +from scrapy.contrib.pbcluster.master.web import ClusterMasterWeb +from scrapy.contrib.pbcluster.crawler.manager import ClusterCrawler \ No newline at end of file diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/crawler/__init__.py b/scrapy/trunk/scrapy/contrib/pbcluster/crawler/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py index 8d66b9799..f75d1a801 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py @@ -1,21 +1,24 @@ +import os + 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 +from scrapy.core.exceptions import NotConfigured class Broker(pb.Referenceable): def __init__(self, crawler, remote): self.__remote = remote self.__crawler = crawler try: - deferred = self.__remote.callRemote("register_crawler", domain, self) + deferred = self.__remote.callRemote("register_crawler", os.getpid(), 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)) + deferred.addCallbacks(callback=lambda x: None, errback=lambda reason: log.msg(reason, log.ERROR)) def remote_stop(self): scrapymanager.stop() @@ -29,5 +32,7 @@ class ClusterCrawler: 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)) + def _set_worker(obj): + self.worker = Broker(self, obj) + d.addCallbacks(callback=_set_worker, errback=lambda reason: log.msg(reason, log.ERROR)) \ No newline at end of file diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py index e5d3255a9..a2da5b7dc 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py @@ -41,7 +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] + del self.procman.crawlers[self.pid] self.procman.update_master(self.domain, "scraped") class ClusterWorker(pb.Root): @@ -53,7 +53,7 @@ class ClusterWorker(pb.Root): self.maxproc = settings.getint('CLUSTER_WORKER_MAXPROC') self.logdir = settings['CLUSTER_LOGDIR'] self.running = {}#a dict domain->ScrapyProcessControl - self.crawlers = {}#a dict domain->scrapy process remote pb connection + self.crawlers = {}#a dict pid->scrapy process remote pb connection self.starttime = datetime.datetime.utcnow() port = settings.getint('CLUSTER_WORKER_PORT') scrapyengine.listenTCP(port, pb.PBServerFactory(self)) @@ -88,7 +88,7 @@ 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)) - d = self.crawler["domain"].callRemote("stop") + d = self.crawlers[proc.pid].callRemote("stop") def _close(): proc.status = "closing" d.addCallbacks(callback=_close, errback=lambda reason: log.msg(reason, log.ERROR)) @@ -123,5 +123,5 @@ class ClusterWorker(pb.Root): 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 + def remote_register_crawler(self, pid, crawler): + self.crawlers[pid] = crawler \ No newline at end of file