diff --git a/scrapy/command/commands/crawl.py b/scrapy/command/commands/crawl.py
index 9c02b08e3..5d7d5ceae 100644
--- a/scrapy/command/commands/crawl.py
+++ b/scrapy/command/commands/crawl.py
@@ -1,5 +1,6 @@
from scrapy import log
from scrapy.command import ScrapyCommand
+from scrapy.core.queue import ExecutionQueue
from scrapy.core.manager import scrapymanager
from scrapy.conf import settings
from scrapy.http import Request
@@ -31,23 +32,25 @@ class Command(ScrapyCommand):
settings.overrides['CRAWLSPIDER_FOLLOW_LINKS'] = False
def run(self, args, opts):
+ q = ExecutionQueue()
urls, names = self._split_urls_and_names(args)
for name in names:
- scrapymanager.crawl_spider_name(name)
+ q.append_spider_name(name)
if opts.spider:
try:
spider = spiders.create(opts.spider)
for url in urls:
- scrapymanager.crawl_url(url, spider)
+ q.append_url(url, spider)
except KeyError:
- log.msg('Could not find spider: %s' % opts.spider, log.ERROR)
+ log.msg('Unable to find spider: %s' % opts.spider, log.ERROR)
else:
for name, urls in self._group_urls_by_spider(urls):
spider = spiders.create(name)
for url in urls:
- scrapymanager.crawl_url(url, spider)
+ q.append_url(url, spider)
+ scrapymanager.queue = q
scrapymanager.start()
def _group_urls_by_spider(self, urls):
diff --git a/scrapy/command/commands/fetch.py b/scrapy/command/commands/fetch.py
index bb2df1400..3ed3674da 100644
--- a/scrapy/command/commands/fetch.py
+++ b/scrapy/command/commands/fetch.py
@@ -28,28 +28,27 @@ class Command(ScrapyCommand):
parser.add_option("--headers", dest="headers", action="store_true", \
help="print response HTTP headers instead of body")
+ def _print_response(self, response, opts):
+ if opts.headers:
+ pprint.pprint(response.headers)
+ else:
+ print response.body
+
def run(self, args, opts):
if len(args) != 1 or not is_url(args[0]):
return False
- responses = [] # to collect downloaded responses
- request = Request(args[0], callback=responses.append, dont_filter=True)
+ cb = lambda x: self._print_response(x, opts)
+ request = Request(args[0], callback=cb, dont_filter=True)
+ spider = None
if opts.spider:
try:
spider = spiders.create(opts.spider)
except KeyError:
log.msg("Could not find spider: %s" % opts.spider, log.ERROR)
- else:
- spider = scrapymanager._create_spider_for_request(request, \
- BaseSpider('default'))
- scrapymanager.crawl_request(request, spider)
+ scrapymanager.configure()
+ scrapymanager.queue.append_request(request, spider, \
+ default_spider=BaseSpider('default'))
scrapymanager.start()
- # display response
- if responses:
- if opts.headers:
- pprint.pprint(responses[0].headers)
- else:
- print responses[0].body
-
diff --git a/scrapy/command/commands/parse.py b/scrapy/command/commands/parse.py
index 9aabe94ff..c1b46edf0 100644
--- a/scrapy/command/commands/parse.py
+++ b/scrapy/command/commands/parse.py
@@ -8,8 +8,6 @@ from scrapy.utils.spider import iterate_spider_output
from scrapy.utils.url import is_url
from scrapy import log
-from collections import defaultdict
-
class Command(ScrapyCommand):
requires_project = True
@@ -75,25 +73,20 @@ class Command(ScrapyCommand):
if not len(args) == 1 or not is_url(args[0]):
return False
- request = Request(args[0])
+ responses = [] # to collect downloaded responses
+ request = Request(args[0], callback=responses.append)
if opts.spider:
try:
spider = spiders.create(opts.spider)
except KeyError:
- log.msg('Could not find spider: %s' % opts.spider, log.ERROR)
+ log.msg('Unable to find spider: %s' % opts.spider, log.ERROR)
return
else:
- spider = scrapymanager._create_spider_for_request(request, \
- log_none=True, log_multiple=True)
+ spider = spiders.create_for_request(request)
- if not spider:
- return
-
- responses = [] # to collect downloaded responses
- request = request.replace(callback=responses.append)
-
- scrapymanager.crawl_request(request, spider)
+ scrapymanager.configure()
+ scrapymanager.queue.append_request(request, spider)
scrapymanager.start()
if not responses:
diff --git a/scrapy/command/commands/runspider.py b/scrapy/command/commands/runspider.py
index 2bfcd87f9..dccd9f465 100644
--- a/scrapy/command/commands/runspider.py
+++ b/scrapy/command/commands/runspider.py
@@ -54,7 +54,7 @@ class Command(ScrapyCommand):
module = _import_file(args[0])
# schedule spider and start engine
- scrapymanager.crawl_spider(module.SPIDER)
+ scrapymanager.queue.append_spider(module.SPIDER)
scrapymanager.start()
if opts.output:
diff --git a/scrapy/command/commands/start.py b/scrapy/command/commands/start.py
index 2c7304787..3c32abcc6 100644
--- a/scrapy/command/commands/start.py
+++ b/scrapy/command/commands/start.py
@@ -1,3 +1,4 @@
+from scrapy.core.queue import KeepAliveExecutionQueue
from scrapy.command import ScrapyCommand
from scrapy.core.manager import scrapymanager
@@ -9,4 +10,6 @@ class Command(ScrapyCommand):
return "Start the Scrapy manager but don't run any spider (idle mode)"
def run(self, args, opts):
- scrapymanager.start(keep_alive=True)
+ q = KeepAliveExecutionQueue()
+ scrapymanager.queue = q
+ scrapymanager.start()
diff --git a/scrapy/contrib/spidermanager.py b/scrapy/contrib/spidermanager.py
index 1c41b625d..263f40f36 100644
--- a/scrapy/contrib/spidermanager.py
+++ b/scrapy/contrib/spidermanager.py
@@ -34,6 +34,27 @@ class TwistedPluginSpiderManager(object):
return [name for name, spider in self._spiders.iteritems()
if url_is_from_spider(request.url, spider)]
+ def create_for_request(self, request, default_spider=None, \
+ log_none=False, log_multiple=False, **spider_kwargs):
+ """Create a spider to handle the given Request.
+
+ This will look for the spiders that can handle the given request (using
+ find_by_request) and return a (new) Spider if (and only if) there is
+ only one Spider able to handle the Request.
+
+ If multiple spiders (or no spider) are found, it will return the
+ default_spider passed. It can optionally log if multiple or no spiders
+ are found.
+ """
+ snames = self.find_by_request(request)
+ if len(snames) == 1:
+ return self.create(snames[0], **spider_kwargs)
+ if len(snames) > 1 and log_multiple:
+ log.msg('More than one spider found for: %s' % request, log.ERROR)
+ if len(snames) == 0 and log_none:
+ log.msg('Unable to find spider for: %s' % request, log.ERROR)
+ return default_spider
+
def list(self):
"""Returns list of spiders available."""
return self._spiders.keys()
diff --git a/scrapy/contrib/webconsole/spiderctl.py b/scrapy/contrib/webconsole/spiderctl.py
index ce3fdffd6..a17c6a3d9 100644
--- a/scrapy/contrib/webconsole/spiderctl.py
+++ b/scrapy/contrib/webconsole/spiderctl.py
@@ -8,7 +8,7 @@ from scrapy.core import signals
from scrapy.core.manager import scrapymanager
from scrapy.core.engine import scrapyengine
from scrapy.spider import spiders
-from scrapy.management.web import banner
+from scrapy.management.web import banner, webconsole_discover_module
from scrapy.conf import settings
class Spiderctl(object):
@@ -20,8 +20,6 @@ class Spiderctl(object):
self.finished = set()
dispatcher.connect(self.spider_opened, signal=signals.spider_opened)
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
-
- from scrapy.management.web import webconsole_discover_module
dispatcher.connect(self.webconsole_discover_module, signal=webconsole_discover_module)
def spider_opened(self, spider):
@@ -35,7 +33,7 @@ class Spiderctl(object):
if wc_request.args:
changes = self.webconsole_control(wc_request)
- self.scheduled = [s.name for s in scrapyengine.spider_scheduler._pending_spiders]
+ self.scheduled = [s[0].name for s in scrapymanager.queue.spider_requests]
self.idle = [d for d in self.enabled_spiders if d not in self.scheduled
and d not in self.running
and d not in self.finished]
@@ -126,8 +124,8 @@ class Spiderctl(object):
if "remove_pending_spiders" in args:
removed = []
for name in args["remove_pending_spiders"]:
- if scrapyengine.spider_scheduler.remove_pending_spider(name):
- removed.append(name)
+ q = scrapymanager.queue
+ q.spider_requests = [x for x in q.spider_requests if x[0].name != name]
if removed:
s += "
"
s += "Removed scheduled spiders:
" % "".join(args["remove_pending_spiders"])
@@ -135,14 +133,14 @@ class Spiderctl(object):
if "add_pending_spiders" in args:
for name in args["add_pending_spiders"]:
if name not in scrapyengine.scheduler.pending_requests:
- scrapymanager.crawl_spider_name(name)
+ scrapymanager.queue.append_spider_name(name)
s += ""
s += "Scheduled spiders:
" % "".join(args["add_pending_spiders"])
s += ""
if "rerun_finished_spiders" in args:
for name in args["rerun_finished_spiders"]:
if name not in scrapyengine.scheduler.pending_requests:
- scrapymanager.crawl_spider_name(name)
+ scrapymanager.queue.append_spider_name(name)
self.finished.remove(name)
s += ""
s += "Re-scheduled finished spiders:
" % "".join(args["rerun_finished_spiders"])
diff --git a/scrapy/core/downloader/manager.py b/scrapy/core/downloader/manager.py
index aec17b05e..16d2fece4 100644
--- a/scrapy/core/downloader/manager.py
+++ b/scrapy/core/downloader/manager.py
@@ -183,10 +183,6 @@ class Downloader(object):
site.cancel_request_calls()
self._process_queue(spider)
- def has_capacity(self):
- """Does the downloader have capacity to handle more spiders"""
- return len(self.sites) < self.concurrent_spiders
-
def is_idle(self):
return not self.sites
diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py
index a26206db6..8bc90ccbe 100644
--- a/scrapy/core/engine.py
+++ b/scrapy/core/engine.py
@@ -6,7 +6,7 @@ For more information see docs/topics/architecture.rst
"""
from time import time
-from twisted.internet import reactor, task, defer
+from twisted.internet import reactor, defer
from twisted.python.failure import Failure
from scrapy.xlib.pydispatch import dispatcher
@@ -27,57 +27,42 @@ class ExecutionEngine(object):
def __init__(self):
self.configured = False
- self.keep_alive = False
self.closing = {} # dict (spider -> reason) of spiders being closed
self.running = False
self.killed = False
self.paused = False
self._next_request_calls = {}
- self._mainloop_task = task.LoopingCall(self._mainloop)
self._crawled_logline = load_object(settings['LOG_FORMATTER_CRAWLED'])
- def configure(self):
+ def configure(self, spider_closed_callback):
"""
Configure execution engine with the given scheduling policy and downloader.
"""
self.scheduler = load_object(settings['SCHEDULER'])()
- self.spider_scheduler = load_object(settings['SPIDER_SCHEDULER'])()
self.downloader = Downloader()
self.scraper = Scraper(self)
self.configured = True
+ self._spider_closed_callback = spider_closed_callback
def start(self):
"""Start the execution engine"""
- if self.running:
- return
+ assert not self.running, "Engine already running"
self.start_time = time()
send_catch_log(signal=signals.engine_started, sender=self.__class__)
- self._mainloop_task.start(5.0, now=True)
- reactor.callWhenRunning(self._mainloop)
self.running = True
def stop(self):
"""Stop the execution engine gracefully"""
- if not self.running:
- return
+ assert self.running, "Engine not running"
self.running = False
- def before_shutdown():
- dfd = self._close_all_spiders()
- return dfd.addBoth(lambda _: self._finish_stopping_engine())
- reactor.addSystemEventTrigger('before', 'shutdown', before_shutdown)
- if self._mainloop_task.running:
- self._mainloop_task.stop()
- try:
- reactor.stop()
- except RuntimeError: # raised if already stopped or in shutdown stage
- pass
+ dfd = self._close_all_spiders()
+ return dfd.addBoth(lambda _: self._finish_stopping_engine())
def kill(self):
"""Forces shutdown without waiting for pending transfers to finish.
stop() must have been called first
"""
- if self.running:
- return
+ assert not self.running, "Call engine.stop() before engine.kill()"
self.killed = True
def pause(self):
@@ -92,12 +77,6 @@ class ExecutionEngine(object):
return self.scheduler.is_idle() and self.downloader.is_idle() and \
self.scraper.is_idle()
- def next_spider(self):
- spider = self.spider_scheduler.next_spider()
- if spider:
- self.open_spider(spider)
- return True
-
def next_request(self, spider, now=False):
"""Scrape the next request for the spider passed.
@@ -163,7 +142,13 @@ class ExecutionEngine(object):
def open_spiders(self):
return self.downloader.sites.keys()
+ def has_capacity(self):
+ """Does the engine have capacity to handle more spiders"""
+ return len(self.downloader.sites) < self.downloader.concurrent_spiders
+
def crawl(self, request, spider):
+ assert spider in self.open_spiders, \
+ "Spider %r not opened when crawling: %s" % (spider.name, request)
if not request.deferred.callbacks:
log.msg("Unable to crawl Request with no callback: %s" % request,
level=log.ERROR, spider=spider)
@@ -180,25 +165,9 @@ class ExecutionEngine(object):
def schedule(self, request, spider):
if spider in self.closing:
raise IgnoreRequest()
- if not self.scheduler.spider_is_open(spider):
- self.scheduler.open_spider(spider)
- if self.spider_is_closed(spider): # scheduler auto-open
- self.spider_scheduler.add_spider(spider)
self.next_request(spider)
return self.scheduler.enqueue_request(spider, request)
- def _mainloop(self):
- """Add more spiders to be scraped if the downloader has the capacity.
-
- If there is nothing else scheduled then stop the execution engine.
- """
- if not self.running or self.paused:
- return
-
- while self.running and self.downloader.has_capacity():
- if not self.next_spider():
- return self._stop_if_idle()
-
def download(self, request, spider):
def _on_success(response):
"""handle the result of a page download"""
@@ -241,14 +210,15 @@ class ExecutionEngine(object):
return dwld
def open_spider(self, spider):
+ assert self.has_capacity(), "No free spider slots when opening %r" % \
+ spider.name
log.msg("Spider opened", spider=spider)
- self.next_request(spider)
-
+ self.scheduler.open_spider(spider)
self.downloader.open_spider(spider)
self.scraper.open_spider(spider)
stats.open_spider(spider)
-
send_catch_log(signals.spider_opened, sender=self.__class__, spider=spider)
+ self.next_request(spider)
def _spider_idle(self, spider):
"""Called when a spider gets idle. This function is called when there
@@ -270,11 +240,6 @@ class ExecutionEngine(object):
if self.spider_is_idle(spider):
self.close_spider(spider, reason='finished')
- def _stop_if_idle(self):
- """Call the stop method if the system has no outstanding tasks. """
- if self.is_idle() and not self.keep_alive:
- self.stop()
-
def close_spider(self, spider, reason='cancelled'):
"""Close (cancel) spider and clear all its outstanding requests"""
if spider not in self.closing:
@@ -316,7 +281,7 @@ class ExecutionEngine(object):
dfd.addErrback(log.err, "Unhandled error on SpiderManager.close_spider()",
spider=spider)
dfd.addBoth(lambda _: log.msg("Spider closed (%s)" % reason, spider=spider))
- reactor.callLater(0, self._mainloop)
+ reactor.callLater(0, self._spider_closed_callback)
return dfd
def _finish_stopping_engine(self):
diff --git a/scrapy/core/manager.py b/scrapy/core/manager.py
index b1dc3754f..d297ed0af 100644
--- a/scrapy/core/manager.py
+++ b/scrapy/core/manager.py
@@ -1,28 +1,26 @@
import signal
-from twisted.internet import reactor
+from twisted.internet import reactor, defer
+from scrapy.core.engine import scrapyengine
+from scrapy.core.queue import ExecutionQueue
from scrapy.extension import extensions
from scrapy import log
-from scrapy.http import Request
-from scrapy.core.engine import scrapyengine
from scrapy.spider import spiders
-from scrapy.utils.misc import arg_to_iter
from scrapy.utils.ossignal import install_shutdown_handlers, signal_names
class ExecutionManager(object):
def __init__(self):
- self.interrupted = False
self.configured = False
self.control_reactor = True
+ self.engine = scrapyengine
- def configure(self, control_reactor=True):
+ def configure(self, control_reactor=True, queue=None):
self.control_reactor = control_reactor
if control_reactor:
install_shutdown_handlers(self._signal_shutdown)
- reactor.addSystemEventTrigger('before', 'shutdown', scrapyengine.stop)
if not log.started:
log.start()
@@ -33,68 +31,54 @@ class ExecutionManager(object):
log.msg("Enabled extensions: %s" % ", ".join(extensions.enabled.iterkeys()),
level=log.DEBUG)
- scrapyengine.configure()
+ self.queue = queue or ExecutionQueue()
+ self.engine.configure(self._spider_closed)
self.configured = True
-
- def crawl_url(self, url, spider=None):
- """Schedule given url for crawling."""
- if spider is None:
- spider = self._create_spider_for_request(Request(url), log_none=True, \
- log_multiple=True)
+
+ @defer.inlineCallbacks
+ def _start_next_spider(self):
+ spider, requests = yield defer.maybeDeferred(self.queue.get_next)
if spider:
- requests = arg_to_iter(spider.make_requests_from_url(url))
- self._crawl_requests(requests, spider)
+ self._start_spider(spider, requests)
+ if self.engine.has_capacity() and not self._nextcall.active():
+ self._nextcall = reactor.callLater(self.queue.polling_delay, \
+ self._start_next_spider)
- def crawl_request(self, request, spider=None):
- """Schedule request for crawling."""
- assert self.configured, "Scrapy Manager not yet configured"
- if spider is None:
- spider = self._create_spider_for_request(request, log_none=True, \
- log_multiple=True)
- if spider:
- scrapyengine.crawl(request, spider)
+ @defer.inlineCallbacks
+ def _start_spider(self, spider, requests):
+ """Don't call this method. Use self.queue to start new spiders"""
+ yield defer.maybeDeferred(self.engine.open_spider, spider)
+ for request in requests:
+ self.engine.crawl(request, spider)
- def crawl_spider_name(self, name):
- """Schedule given spider by name for crawling."""
- try:
- spider = spiders.create(name)
- except KeyError:
- log.msg('Could not find spider: %s' % name, log.ERROR)
- else:
- self.crawl_spider(spider)
+ @defer.inlineCallbacks
+ def _spider_closed(self):
+ if not self.engine.open_spiders:
+ is_finished = yield defer.maybeDeferred(self.queue.is_finished)
+ if is_finished:
+ self.stop()
+ return
+ if self.engine.has_capacity():
+ self._start_next_spider()
- def crawl_spider(self, spider):
- """Schedule spider for crawling."""
- requests = spider.start_requests()
- self._crawl_requests(requests, spider)
-
- def _crawl_requests(self, requests, spider):
- """Shortcut to schedule a list of requests"""
- for req in requests:
- self.crawl_request(req, spider)
-
- def start(self, keep_alive=False):
- """Start the scrapy server, without scheduling any domains"""
- scrapyengine.keep_alive = keep_alive
- scrapyengine.start()
+ @defer.inlineCallbacks
+ def start(self):
+ yield defer.maybeDeferred(self.engine.start)
+ self._nextcall = reactor.callLater(0, self._spider_closed)
+ reactor.addSystemEventTrigger('before', 'shutdown', self.stop)
if self.control_reactor:
reactor.run(installSignalHandlers=False)
+ @defer.inlineCallbacks
def stop(self):
- """Stop the scrapy server, shutting down the execution engine"""
- self.interrupted = True
- scrapyengine.stop()
-
- def _create_spider_for_request(self, request, default=None, log_none=False, \
- log_multiple=False):
- spider_names = spiders.find_by_request(request)
- if len(spider_names) == 1:
- return spiders.create(spider_names[0])
- if len(spider_names) > 1 and log_multiple:
- log.msg('More than one spider found for: %s' % request, log.ERROR)
- if len(spider_names) == 0 and log_none:
- log.msg('Could not find spider for: %s' % request, log.ERROR)
- return default
+ if self._nextcall.active():
+ self._nextcall.cancel()
+ if self.engine.running:
+ yield defer.maybeDeferred(self.engine.stop)
+ try:
+ reactor.stop()
+ except RuntimeError: # raised if already stopped or in shutdown stage
+ pass
def _signal_shutdown(self, signum, _):
signame = signal_names[signum]
@@ -108,7 +92,7 @@ class ExecutionManager(object):
log.msg('Received %s twice, forcing unclean shutdown' % signame, \
level=log.INFO)
log.log_level = log.SILENT # disable logging of confusing tracebacks
- reactor.callFromThread(scrapyengine.kill)
+ reactor.callFromThread(self.engine.kill)
install_shutdown_handlers(signal.SIG_IGN)
scrapymanager = ExecutionManager()
diff --git a/scrapy/core/queue.py b/scrapy/core/queue.py
new file mode 100644
index 000000000..f8dd340df
--- /dev/null
+++ b/scrapy/core/queue.py
@@ -0,0 +1,87 @@
+from twisted.internet import defer
+
+from scrapy.http import Request
+from scrapy.utils.misc import arg_to_iter
+from scrapy import log
+from scrapy.spider import spiders
+
+
+class ExecutionQueue(object):
+
+ polling_delay = 5
+
+ def __init__(self, _spiders=spiders):
+ self.spider_requests = []
+ self._spiders = _spiders
+
+ def _append_next(self):
+ """Called when there are no more itemsl left in self.spider_requests.
+ This method is meant to be overriden in subclasses to add new (spider,
+ requests) tuples to self.spider_requests. It can return a Deferred.
+ """
+ pass
+
+ def get_next(self):
+ """Return a tuple (spider, requests) containing a list of Requests and
+ the Spider which will be used to crawl those Requests. If there aren't
+ any more spiders to crawl it must return (None, []).
+
+ This method can return a deferred.
+ """
+ if self.spider_requests:
+ return self._get_next_now()
+ d = defer.maybeDeferred(self._append_next)
+ d.addCallback(lambda _: self._get_next_now())
+ return d
+
+ def _get_next_now(self):
+ try:
+ return self.spider_requests.pop(0)
+ except IndexError:
+ return (None, [])
+
+ def is_finished(self):
+ """Return True if the queue is empty and there won't be any more
+ spiders to crawl (this is for one-shot runs). If it returns ``False``
+ Scrapy will keep polling this queue for new requests to scrape
+ """
+ return not bool(self.spider_requests)
+
+ def append_spider(self, spider):
+ """Append a Spider to crawl"""
+ requests = spider.start_requests()
+ self.spider_requests.append((spider, requests))
+
+ def append_request(self, request, spider=None, **kwargs):
+ if spider is None:
+ spider = self._spiders.create_for_request(request, **kwargs)
+ if spider:
+ self.spider_requests.append((spider, [request]))
+
+ def append_url(self, url, spider=None, **kwargs):
+ """Append a URL to crawl with the given spider. If the spider is not
+ given, a spider will be looked up based on the URL
+ """
+ if spider is None:
+ spider = self._spiders.create_for_request(Request(url), **kwargs)
+ if spider:
+ requests = arg_to_iter(spider.make_requests_from_url(url))
+ self.spider_requests.append((spider, requests))
+
+ def append_spider_name(self, spider_name, **spider_kwargs):
+ """Append a spider to crawl given its name and optional arguments,
+ which are used to instantiate it. The SpiderManager is used to lookup
+ the spider
+ """
+ try:
+ spider = self._spiders.create(spider_name, **spider_kwargs)
+ except KeyError:
+ log.msg('Unable to find spider: %s' % spider_name, log.ERROR)
+ else:
+ self.append_spider(spider)
+
+
+class KeepAliveExecutionQueue(ExecutionQueue):
+
+ def is_finished(self):
+ return False
diff --git a/scrapy/shell.py b/scrapy/shell.py
index 96dc51270..97af70869 100644
--- a/scrapy/shell.py
+++ b/scrapy/shell.py
@@ -18,6 +18,7 @@ from scrapy.utils.response import open_in_browser
from scrapy.conf import settings
from scrapy.core.manager import scrapymanager
from scrapy.core.engine import scrapyengine
+from scrapy.core.queue import KeepAliveExecutionQueue
from scrapy.http import Request, TextResponse
def relevant_var(varname):
@@ -54,10 +55,11 @@ class Shell(object):
url = parse_url(request_or_url)
request = Request(url)
- spider = scrapymanager._create_spider_for_request(request, \
- BaseSpider('default'), log_multiple=True)
+ spider = spiders.create_for_request(request, BaseSpider('default'), \
+ log_multiple=True)
print "Fetching %s..." % request
+ scrapyengine.open_spider(spider)
response = threads.blockingCallFromThread(reactor, scrapyengine.schedule, \
request, spider)
if response:
@@ -108,7 +110,8 @@ class Shell(object):
signal.signal(signal.SIGINT, signal.SIG_IGN)
reactor.callInThread(self._console_thread, url)
- scrapymanager.start(keep_alive=True)
+ scrapymanager.queue = KeepAliveExecutionQueue()
+ scrapymanager.start()
def inspect_response(self, response):
print
diff --git a/scrapy/tests/test_core_queue.py b/scrapy/tests/test_core_queue.py
new file mode 100644
index 000000000..2d41910ad
--- /dev/null
+++ b/scrapy/tests/test_core_queue.py
@@ -0,0 +1,105 @@
+import unittest
+
+from scrapy.core.queue import ExecutionQueue, KeepAliveExecutionQueue
+from scrapy.spider import BaseSpider
+from scrapy.http import Request
+
+class TestSpider(BaseSpider):
+
+ name = "default"
+
+ def start_requests(self):
+ return [Request("http://www.example.com/1"), \
+ Request("http://www.example.com/2")]
+
+ def make_requests_from_url(self, url):
+ return [Request(url + "/make1"), Request(url + "/make2")]
+
+
+class TestSpiderManager(object):
+
+ def create_for_request(self, request, **kwargs):
+ return TestSpider('create_for_request', **kwargs)
+
+ def create(self, spider_name, **spider_kwargs):
+ return TestSpider(spider_name, **spider_kwargs)
+
+
+class ExecutionQueueTest(unittest.TestCase):
+
+ queue_class = ExecutionQueue
+
+ def setUp(self):
+ self.queue = self.queue_class(_spiders=TestSpiderManager())
+ self.spider = TestSpider()
+ self.request = Request('about:none')
+
+ def tearDown(self):
+ del self.queue, self.spider, self.request
+
+ def test_is_finished(self):
+ self.assert_(self.queue.is_finished())
+ self.queue.append_request(self.request, self.spider)
+ self.assert_(not self.queue.is_finished())
+
+ def test_append_spider(self):
+ spider = TestSpider()
+ self.queue.append_spider(spider)
+ self.assert_(self.queue.spider_requests[0][0] is spider)
+ self._assert_request_urls(self.queue.spider_requests[0][1],
+ ["http://www.example.com/1", "http://www.example.com/2"])
+
+ def test_append_request1(self):
+ spider = TestSpider()
+ request = Request('about:blank')
+ self.queue.append_request(request, spider=spider)
+ self.assert_(self.queue.spider_requests[0][0] is spider)
+ self.assert_(self.queue.spider_requests[0][1][0] is request)
+
+ def test_append_request2(self):
+ request = Request('about:blank')
+ self.queue.append_request(request, arg='123')
+ spider = self.queue.spider_requests[0][0]
+ self.assert_(spider.name == 'create_for_request')
+ self.assert_(spider.arg == '123')
+
+ def test_append_url(self):
+ spider = TestSpider()
+ url = 'http://www.example.com/asd'
+ self.queue.append_url(url, spider=spider)
+ self.assert_(self.queue.spider_requests[0][0] is spider)
+ self._assert_request_urls(self.queue.spider_requests[0][1], \
+ ['http://www.example.com/asd/make1', 'http://www.example.com/asd/make2'])
+
+ def test_append_url2(self):
+ url = 'http://www.example.com/asd'
+ self.queue.append_url(url, arg='123')
+ self._assert_request_urls(self.queue.spider_requests[0][1], \
+ ['http://www.example.com/asd/make1', 'http://www.example.com/asd/make2'])
+ spider = self.queue.spider_requests[0][0]
+ self.assert_(spider.name == 'create_for_request')
+ self.assert_(spider.arg == '123')
+
+ def test_append_spider_name(self):
+ self.queue.append_spider_name('test123', arg='123')
+ spider = self.queue.spider_requests[0][0]
+ self.assert_(spider.name == 'test123')
+ self.assert_(spider.arg == '123')
+
+ def _assert_request_urls(self, requests, urls):
+ assert all(isinstance(x, Request) for x in requests)
+ self.assertEqual([x.url for x in requests], urls)
+
+class KeepAliveExecutionQueueTest(ExecutionQueueTest):
+
+ queue_class = KeepAliveExecutionQueue
+
+ def test_is_finished(self):
+ self.assert_(not self.queue.is_finished())
+ self.queue.append_request(self.request, self.spider)
+ self.assert_(not self.queue.is_finished())
+
+
+if __name__ == "__main__":
+ unittest.main()
+
diff --git a/scrapy/tests/test_engine.py b/scrapy/tests/test_engine.py
index 2d7b5950c..8fabd2913 100644
--- a/scrapy/tests/test_engine.py
+++ b/scrapy/tests/test_engine.py
@@ -97,7 +97,7 @@ class CrawlingSession(object):
dispatcher.connect(self.response_downloaded, signals.response_downloaded)
scrapymanager.configure()
- scrapymanager.crawl_spider(self.spider)
+ scrapymanager.queue.append_spider(self.spider)
scrapymanager.start()
self.port.stopListening()
self.wasrun = True
diff --git a/scrapy/utils/engine.py b/scrapy/utils/engine.py
index b1250c71e..e573730ab 100644
--- a/scrapy/utils/engine.py
+++ b/scrapy/utils/engine.py
@@ -12,11 +12,11 @@ def get_engine_status(engine=None):
global_tests = [
"time()-engine.start_time",
"engine.is_idle()",
+ "engine.has_capacity()",
"engine.scheduler.is_idle()",
"len(engine.scheduler.pending_requests)",
"engine.downloader.is_idle()",
"len(engine.downloader.sites)",
- "engine.downloader.has_capacity()",
"engine.scraper.is_idle()",
"len(engine.scraper.sites)",
]