diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py index 85299a61d..19ad1a96b 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py @@ -71,12 +71,12 @@ class Node: def _run_callback(status): if status['callresponse'][0] == 1: #slots are complete. Reschedule in master. This is a security issue because could happen that the slots were completed since last status update by another cluster (thinking at future with full-distributed worker-master clusters) - self.master.schedule(pending['domain'], pending['settings'], pending['priority']) + self.master.schedule(pending['domain'], pending['settings'], pending['priority'], pending["env"]) log.msg("Domain %s rescheduled: no proc space in node." % pending['domain'], log.WARNING) self._set_status(status) try: - deferred = self.__remote.callRemote("run", pending["domain"], pending["settings"]) + deferred = self.__remote.callRemote("run", pending["domain"], pending["settings"], pending["env"]) except pb.DeadReferenceError: self._set_status(None) log.msg("Lost connection to node %s." % (self.name), log.ERROR) @@ -132,7 +132,7 @@ class ClusterMaster(object): def remove_node(self, nodename): raise NotImplemented - def schedule(self, domains, spider_settings=None, priority=PRIORITY_NORMAL): + def schedule(self, domains, spider_settings=None, env=None, priority=PRIORITY_NORMAL): i = 0 for p in self.pending: if p['priority'] <= priority: @@ -140,7 +140,7 @@ class ClusterMaster(object): else: break for domain in domains: - self.pending.insert(i, {'domain': domain, 'settings': spider_settings, 'priority': priority}) + self.pending.insert(i, {'domain': domain, 'settings': spider_settings, 'env': env, 'priority': priority}) self.update_nodes() def stop(self, domains): diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py index 1a2c9aae1..be384b18b 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py @@ -62,6 +62,8 @@ class ClusterMasterWeb(ClusterMaster): sep = "\r" domains = args["schedule"] priority = eval(args.get("priority", ["PRIORITY_NORMAL"])[0]) + + #spider settings slist = args.get("settings", [""])[0].split(sep) spider_settings = {} for s in slist: @@ -71,7 +73,19 @@ class ClusterMasterWeb(ClusterMaster): pass else: spider_settings[k] = v - self.schedule(domains, spider_settings, priority) + + #other environment settings + envlist = args.get("env", [""])[0].split(sep) + env = {} + for e in envlist: + try: + k, v = e.strip().split("=") + except ValueError: + pass + else: + env[k] = v + + self.schedule(domains, spider_settings, env, priority) if ws: return self.ws_status(wc_request) @@ -157,11 +171,20 @@ class ClusterMasterWeb(ClusterMaster): s += "" % (p, pname) s += "\n" s += "
\n" + + #spider settings s += "Spider settings:
\n" s += "\n" + s += "
\n" + + #other environment settings + s += "Other environment settings:
\n" + s += "\n" + s += "

\n" s += "\n" @@ -210,12 +233,7 @@ class ClusterMasterWeb(ClusterMaster): wc_request.setHeader('content-type', 'text/plain') status = {} nodes_status = {} - running = [] - for d, n in self.nodes.iteritems(): - nodes_status[d] = n.status_as_dict - running.extend(n.running) status["nodes"] = nodes_status status["pending"] = self.pending - status["running"] = running content = serialize(status, format) return content \ 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 1641234cf..1472832ed 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py @@ -12,26 +12,26 @@ from scrapy.conf import settings from scrapy.core.engine import scrapyengine class ScrapyProcessProtocol(protocol.ProcessProtocol): - def __init__(self, procman, domain, logfile=None, spider_settings=None): + def __init__(self, procman, domain, logfile=None, spider_settings=None, env=None): self.procman = procman self.domain = domain self.logfile = logfile self.start_time = datetime.datetime.utcnow() self.status = "starting" self.pid = -1 - - env = {'SCRAPY_LOGFILE': self.logfile, 'SCRAPY_CLUSTER_WORKER_ENABLED': '0', 'SCRAPY_WEBCONSOLE_ENABLED': '0'} + self.env_original = env or {} + self.env = self.env_original.copy() #We conserve original setting format for info purposes (avoid lots of unnecesary "SCRAPY_") - self.settings = spider_settings or {} - for k in self.settings: - env["SCRAPY_%s" % k] = self.settings[k] - self.env = env + self.scrapy_settings = spider_settings or {} + self.scrapy_settings.update({'LOGFILE': self.logfile, 'CLUSTER_WORKER_ENABLED': '0', 'WEBCONSOLE_ENABLED': '0'}) + for k in self.scrapy_settings: + self.env["SCRAPY_%s" % k] = self.scrapy_settings[k] def __str__(self): return "" % (self.domain, self.pid, self.status) def as_dict(self): - return {"domain": self.domain, "pid": self.pid, "status": self.status, "settings": self.settings, "logfile": self.logfile, "starttime": self.start_time} + return {"domain": self.domain, "pid": self.pid, "status": self.status, "settings": self.scrapy_settings, "logfile": self.logfile, "starttime": self.start_time, "env": self.env_original} def connectionMade(self): self.pid = self.transport.pid @@ -81,13 +81,13 @@ class ClusterWorker(pb.Root): status["callresponse"] = (rcode, rstring) if rstring else (0, "Status Response.") return status - def remote_run(self, domain, spider_settings=None): + def remote_run(self, domain, spider_settings=None, env=None): """Spawn process to run the given domain.""" if len(self.running) < self.maxproc: logfile = os.path.join(self.logdir, domain, time.strftime("%FT%T.log")) if not os.path.exists(os.path.dirname(logfile)): os.makedirs(os.path.dirname(logfile)) - scrapy_proc = ScrapyProcessProtocol(self, domain, logfile, spider_settings) + scrapy_proc = ScrapyProcessProtocol(self, domain, logfile, spider_settings, env) args = [sys.executable, sys.argv[0], 'crawl', domain] proc = reactor.spawnProcess(scrapy_proc, sys.executable, args=args, env=scrapy_proc.env) self.running[domain] = scrapy_proc