diff --git a/scrapy/trunk/scrapy/contrib/cluster/__init__.py b/scrapy/trunk/scrapy/contrib/cluster/__init__.py deleted file mode 100644 index 5c544146d..000000000 --- a/scrapy/trunk/scrapy/contrib/cluster/__init__.py +++ /dev/null @@ -1,2 +0,0 @@ -from scrapy.contrib.cluster.master import ClusterMaster, ClusterMasterWeb -from scrapy.contrib.cluster.worker import ClusterWorker, ClusterWorkerWeb diff --git a/scrapy/trunk/scrapy/contrib/cluster/master/__init__.py b/scrapy/trunk/scrapy/contrib/cluster/master/__init__.py deleted file mode 100644 index fdb4d71b9..000000000 --- a/scrapy/trunk/scrapy/contrib/cluster/master/__init__.py +++ /dev/null @@ -1,2 +0,0 @@ -from scrapy.contrib.cluster.master.manager import ClusterMaster -from scrapy.contrib.cluster.master.web import ClusterMasterWeb diff --git a/scrapy/trunk/scrapy/contrib/cluster/master/manager.py b/scrapy/trunk/scrapy/contrib/cluster/master/manager.py deleted file mode 100644 index 32a5f9239..000000000 --- a/scrapy/trunk/scrapy/contrib/cluster/master/manager.py +++ /dev/null @@ -1,189 +0,0 @@ -import urlparse -import urllib -import bisect - -from pydispatch import dispatcher -from twisted.web.client import getPage - -from scrapy.core import log, signals -from scrapy.core.engine import scrapyengine -from scrapy.core.exceptions import NotConfigured -from scrapy.utils.serialization import unserialize -from scrapy.conf import settings - -class ClusterNode(object): - def __init__(self, name, url): - self.name = name - self.url = url - self.maxproc = 0 - self.running = {} - self.pending = [] - self.loadavg = (0.0, 0.0, 0.0) - self.status = "down" # down/crawling/idle/error - self.available = False - - self.wsurl = urlparse.urljoin(self.url, "cluster_worker/ws/?format=json") - - def update(self): - d = getPage(self.wsurl) - d.addCallbacks(self._cbUpdate, self._ebUpdate) - - def schedule(self, domains): - args = [("schedule", domain) for domain in domains] - self._wsRequest(args) - - def stop(self, domains): - args = [("stop", domain) for domain in domains] - self._wsRequest(args) - - def remove(self, domains): - args = [("remove", domain) for domain in domains] - self._wsRequest(args) - - def _wsRequest(self, args): - d = getPage("%s?%s" % (self.wsurl, urllib.urlencode(args))) - d.addCallbacks(self._cbUpdate, self._ebUpdate) - - def _cbUpdate(self, jsonstatus): - try: - pmstatus = unserialize(jsonstatus, 'json') - self.maxproc = int(pmstatus.get('maxproc', None)) - self.running = pmstatus.get('running') or {} - self.pending = pmstatus.get('pending') or [] - self.loadavg = tuple(pmstatus.get('loadavg', (0.0, 0.0, 0.0))) - self.status = "crawling" if self.running else "idle" - self.available = True - except ValueError: - self.status = "error" - self.available = False - - def _ebUpdate(self, err): - self.status = "down" - self.available = False - - -class ClusterMaster(object): - - def __init__(self): - if not settings.getbool('CLUSTER_MASTER_ENABLED'): - raise NotConfigured - self.nodes = {} - - dispatcher.connect(self._engine_started, signal=signals.engine_started) - - def load_nodes(self): - """Loads nodes from the CLUSTER_MASTER_NODES setting""" - for name, url in settings.get('CLUSTER_MASTER_NODES', {}).iteritems(): - self.add_node(name, url) - - def update_nodes(self): - for node in self.nodes.itervalues(): - node.update() - - def add_node(self, name, url): - """Add node given its node""" - if name not in self.nodes: - node = ClusterNode(name, url) - self.nodes[name] = node - node.update() - - def remove_node(self, nodename): - raise NotImplemented - - def schedule(self, domains, nodename=None): - if nodename: - node = self.nodes[nodename] - node.schedule(domains) - else: - self._dispatch_domains(domains) - - def stop(self, domains): - to_stop = {} - for domain in domains: - node = self.running.get(domain, None) - if node: - if node.name not in to_stop: - to_stop[node.name] = [] - to_stop[node.name].append(domain) - - for nodename, domains in to_stop.iteritems(): - self.nodes[nodename].stop(domains) - - def remove(self, domains): - to_remove = {} - for domain in domains: - node = self.pending.get(domain, None) - if node: - if node.name not in to_remove: - to_remove[node.name] = [] - to_remove[node.name].append(domain) - - for nodename, domains in to_remove.iteritems(): - self.nodes[nodename].remove(domains) - - def discard(self, domains): - """Stop and remove all running and pending instances of the given - domains""" - self.remove(domains) - self.stop(domains) - - @property - def running(self): - """Return dict of running domains as domain -> node""" - d = {} - for node in self.nodes.itervalues(): - for domain in node.running.iterkeys(): - d[domain] = node - return d - - @property - def pending(self): - """Return dict of pending domains as domain -> node""" - d = {} - for node in self.nodes.itervalues(): - for domain in node.pending: - d[domain] = node - return d - - @property - def available_nodes(self): - return (node for node in self.nodes.itervalues() if node.available) - - def _dispatch_domains(self, domains): - """Schedule the given domains in the availables nodes as good as - possible. The algorithm follows the next rules (in order): - - 1. search for nodes with available capacity(running < maxproc) and (if - any) schedules the domains there - - 2. if there isn't any node with available capacity it schedules the - domain in the node with the smallest number of pending spiders - """ - - to_schedule = {} # domains to schedule per node - pending_node = [] # list of #pending, node - - for node in self.available_nodes: - capacity = node.maxproc - len(node.running) - bisect.insort(pending_node, (len(node.pending), node)) - - to_schedule[node.name] = [] - while domains and capacity > 0: - to_schedule[node.name].append(domains.pop(0)) - capacity -= 1 - if not domains: - break - - for domain in domains: - pending, node = pending_node.pop(0) - to_schedule[node.name].append(domain) - bisect.insort(pending_node, (pending+1, node)) - - for nodename, domains in to_schedule.iteritems(): - if domains: - self.nodes[nodename].schedule(domains) - - def _engine_started(self): - self.load_nodes() - scrapyengine.addtask(self.update_nodes, settings.getint('CLUSTER_MASTER_POLL_INTERVAL')) - diff --git a/scrapy/trunk/scrapy/contrib/cluster/master/web.py b/scrapy/trunk/scrapy/contrib/cluster/master/web.py deleted file mode 100644 index 984d1ba8a..000000000 --- a/scrapy/trunk/scrapy/contrib/cluster/master/web.py +++ /dev/null @@ -1,165 +0,0 @@ -import datetime - -from pydispatch import dispatcher - -from scrapy.spider import spiders -from scrapy.management.web import banner, webconsole_discover_module -from scrapy.utils.serialization import parse_jsondatetime -from scrapy.contrib.cluster.master import ClusterMaster - -class ClusterMasterWeb(ClusterMaster): - webconsole_id = 'cluster_master' - webconsole_name = 'Cluster master' - - def __init__(self): - ClusterMaster.__init__(self) - - dispatcher.connect(self.webconsole_discover_module, signal=webconsole_discover_module) - - def webconsole_render(self, wc_request): - changes = "" - if wc_request.path == '/cluster/nodes/': - return self.render_nodes(wc_request) - elif wc_request.path == '/cluster/domains/': - return self.render_domains(wc_request) - elif wc_request.args: - changes = self.webconsole_control(wc_request) - - s = self.render_header() - - s += "

Home

\n" - - s += "\n" - s += "\n" - for node in self.nodes.itervalues(): - #chkbox = "" % domain if node.status in ["up", "idle"] else " " - nodelink = "%s" % (node.name, node.name) - chkbox = " " - loadavg = "%.2f %.2f %.2f" % node.loadavg - s += "\n" % \ - (chkbox, nodelink, node.status, len(node.running), node.maxproc, len(node.pending), loadavg) - s += "
 NameStatusRunningPendingLoad.avg
%s%s%s%d/%d%d%s
\n" - - s += "\n" - s += "\n" - - return str(s) - - def webconsole_control(self, wc_request): - args = wc_request.args - - if "updatenodes" in args: - self.update_nodes() - - if "schedule" in args: - node = args["node"][0] if "node" in args else None - self.schedule(args["schedule"], nodename=node) - - if "stop" in args: - self.stop(args["stop"]) - - if "remove" in args: - self.remove(args["remove"]) - - return "" - - def render_nodes(self, wc_request): - if wc_request.args: - self.webconsole_control(wc_request) - - now = datetime.datetime.utcnow() - - s = self.render_header() - for node in self.nodes.itervalues(): - s += "

%s

\n" % (node.name, node.name) - - s += "

Running domains

\n" - if node.running: - s += "
\n" - s += "\n" - s += "\n" - for domain, proc in node.running.iteritems(): - chkbox = "" % domain if proc['status'] == "running" else " " - start_time = parse_jsondatetime(proc.get('start_time', None)) - elapsed = now - start_time if start_time else None - s += "\n" % \ - (chkbox, proc['pid'], domain, proc['status'], elapsed, proc['logfile']) - s += "
 PIDDomainStatusRunning timeLog file
%s%s%s%s%s%s
\n" - s += "\n" % node.name - s += "

\n" % node.name - s += "
\n" - else: - s += "

No running domains on %s

\n" % node.name - - # pending domains - s += "

Pending domains

\n" - if node.pending: - s += "
\n" - s += "\n" - s += "\n" % node.name - s += "

\n" % node.name - s += "
\n" - else: - s += "

No pending domains on %s

\n" % node.name - - return str(s) - - def render_domains(self, wc_request): - if wc_request.args: - self.webconsole_control(wc_request) - - enabled_domains = set(spiders.asdict(include_disabled=False).keys()) - inactive_domains = enabled_domains - set(self.running.keys() + self.pending.keys()) - - s = self.render_header() - - s += "

Schedule domains

\n" - - s += "Inactive domains (not running or pending)
" - s += "
\n" - s += "\n" - s += "
\n" - s += "Node (only available nodes shown):
\n" - s += "\n" - s += "

\n" - s += "
\n" - - s += "

Domains

\n" - - s += "\n" - s += "\n" - s += self._domains_table(self.running, 'running') - s += self._domains_table(self.pending, 'pending') - s += "
DomainStatusNode
\n" - - return str(s) - - def render_header(self): - s = banner(self) - s += "

Nav: " - s += "Home | " - s += "Domains | " - s += "Nodes (update)" - s += "

" - return s - - def _domains_table(self, dict_, status): - s = "" - for domain, node in dict_.iteritems(): - s += "%s%s%s\n" % (domain, status, node.name) - return s - - def webconsole_discover_module(self): - return self diff --git a/scrapy/trunk/scrapy/contrib/cluster/worker/__init__.py b/scrapy/trunk/scrapy/contrib/cluster/worker/__init__.py deleted file mode 100644 index 7b97b63ed..000000000 --- a/scrapy/trunk/scrapy/contrib/cluster/worker/__init__.py +++ /dev/null @@ -1,2 +0,0 @@ -from scrapy.contrib.cluster.worker.manager import ClusterWorker -from scrapy.contrib.cluster.worker.web import ClusterWorkerWeb diff --git a/scrapy/trunk/scrapy/contrib/cluster/worker/manager.py b/scrapy/trunk/scrapy/contrib/cluster/worker/manager.py deleted file mode 100644 index 8914c9e0a..000000000 --- a/scrapy/trunk/scrapy/contrib/cluster/worker/manager.py +++ /dev/null @@ -1,96 +0,0 @@ -import sys -import os -import time -import datetime - -from twisted.internet import protocol, reactor - -from scrapy.core import log -from scrapy.core.exceptions import NotConfigured -from scrapy.conf import settings - -class ScrapyProcessProtocol(protocol.ProcessProtocol): - def __init__(self, procman, domain, logfile=None): - self.procman = procman - self.domain = domain - self.logfile = logfile - self.start_time = datetime.datetime.utcnow() - self.status = "starting" - self.pid = -1 - - def __str__(self): - return "" % (self.domain, self.pid, self.status) - - def connectionMade(self): - self.pid = self.transport.pid - log.msg("ClusterWorker: started domain=%s, pid=%d, log=%s" % (self.domain, self.pid, self.logfile)) - self.transport.closeStdin() - self.status = "running" - - def processEnded(self, status_object): - log.msg("ClusterWorker: finished domain=%s, pid=%d, log=%s" % (self.domain, self.pid, self.logfile)) - del self.procman.running[self.domain] - self.procman.next_pending() - - -class ClusterWorker(object): - - def __init__(self): - if not settings.getbool('CLUSTER_WORKER_ENABLED'): - raise NotConfigured - - self.maxproc = settings.getint('CLUSTER_WORKER_MAXPROC') - self.logdir = settings['CLUSTER_WORKER_LOGDIR'] - self.running = {} - self.pending = [] - - def schedule(self, domain): - """Schedule new domain to be crawled in a separate processes""" - - if len(self.running) < self.maxproc and domain not in self.running: - self._run(domain) - else: - self.pending.append(domain) - - def stop(self, domain): - """Stop running domain. For removing pending (not yet started) domains - use remove() instead""" - - 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" - - def remove(self, domain): - """Remove all scheduled instances of the given domain (if it hasn't - started yet). Otherwise use stop()""" - - while domain in self.pending: - self.pending.remove(domain) - - def next_pending(self): - """Run the next domain in the pending list, which is not already running""" - - if len(self.running) >= self.maxproc: - return - for domain in self.pending: - if domain not in self.running: - self._run(domain) - self.pending.remove(domain) - return - - def _run(self, domain): - """Spawn process to run the given domain. Don't call this method - directly. Instead use schedule().""" - - 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) - env = {'SCRAPY_LOGFILE': logfile} - args = [sys.executable, sys.argv[0], 'crawl', domain] - proc = reactor.spawnProcess(scrapy_proc, sys.executable, args=args, env=env) - self.running[domain] = scrapy_proc - - diff --git a/scrapy/trunk/scrapy/contrib/cluster/worker/web.py b/scrapy/trunk/scrapy/contrib/cluster/worker/web.py deleted file mode 100644 index 889d90a3f..000000000 --- a/scrapy/trunk/scrapy/contrib/cluster/worker/web.py +++ /dev/null @@ -1,127 +0,0 @@ -import os -import datetime - -from pydispatch import dispatcher - -from scrapy.spider import spiders -from scrapy.management.web import banner, webconsole_discover_module -from scrapy.utils.serialization import serialize -from scrapy.contrib.cluster.worker import ClusterWorker - -class ClusterWorkerWeb(ClusterWorker): - webconsole_id = 'cluster_worker' - webconsole_name = 'Cluster worker' - - def __init__(self): - ClusterWorker.__init__(self) - - dispatcher.connect(self.webconsole_discover_module, signal=webconsole_discover_module) - - def webconsole_render(self, wc_request): - changes = "" - if wc_request.path == '/cluster_worker/ws/': - return self.webconsole_control(wc_request, ws=True) - elif wc_request.args: - changes = self.webconsole_control(wc_request) - - now = datetime.datetime.utcnow() - - s = banner(self) - - # running processes - s += "

Running processes

\n" - if self.running: - s += "
\n" - s += "\n" - s += "\n" - for domain, proc in self.running.iteritems(): - chkbox = "" % domain if proc.status == "running" else " " - elapsed = now - proc.start_time - s += "\n" % \ - (chkbox, proc.pid, domain, proc.status, proc.logfile, elapsed) - s += "
 PIDDomainStatusLog fileRunning time
%s%s%s%s%s%s
\n" - s += "

\n" - s += "
\n" - else: - s += "

No running processes

\n" - - # pending domains - s += "

Pending domains

\n" - if self.pending: - s += "
\n" - s += "\n" - s += "

\n" - s += "
\n" - else: - s += "

No pending domains

\n" - - # schedule domains - enabled_domains = spiders.asdict(include_disabled=False).keys() - s += "

Schedule domains

\n" - s += "
\n" - s += "\n" - s += "

\n" - s += "
\n" - - s += changes - - s += "\n" - s += "\n" - - return s - - def webconsole_control(self, wc_request, ws=False): - args = wc_request.args - - if "schedule" in args: - for domain in args["schedule"]: - self.schedule(domain) - if ws: - return self.ws_status(wc_request) - else: - return "

Scheduled domains:

" % "
  • ".join(args["schedule"]) + "

    \n" - - if "stop" in args: - for domain in args["stop"]: - self.stop(domain) - if ws: - return self.ws_status(wc_request) - else: - return "

    Stopped running processes:

    " % "
  • ".join(args["stop"]) + "

    \n" - - if "remove" in args: - for domain in args["remove"]: - self.remove(domain) - if ws: - return self.ws_status(wc_request) - else: - return "

    Removed pending domains:

    " % "
  • ".join(args["remove"]) + "

    \n" - - if ws: - return self.ws_status(wc_request) - else: - return "" - - def ws_status(self, wc_request): - format = wc_request.args['format'][0] if 'format' in wc_request.args else 'json' - wc_request.setHeader('content-type', 'text/plain') - exported_proc_attrs = ['pid', 'status', 'start_time', 'logfile'] - d = {'maxproc': self.maxproc, 'running': {}, 'pending': []} - for domain, proc in self.running.iteritems(): - d2 = {} - for a in exported_proc_attrs: - d2[a] = getattr(proc, a) - d['running'][domain] = d2 - d['pending'] = self.pending - d['loadavg'] = os.getloadavg() - content = serialize(d, format) - return content - - def webconsole_discover_module(self): - return self