mirror of https://github.com/scrapy/scrapy.git
remove scrapyd, it was migrated to its own repository
This commit is contained in:
parent
33ca295129
commit
5db45b3825
|
|
@ -9,7 +9,7 @@ env:
|
|||
install:
|
||||
- pip install --use-mirrors -r .travis/requirements-$BUILDENV.txt
|
||||
script:
|
||||
- trial scrapy scrapyd
|
||||
- trial scrapy
|
||||
notifications:
|
||||
- irc:
|
||||
- "irc.freenode.org#scrapy"
|
||||
|
|
|
|||
|
|
@ -4,8 +4,6 @@ include INSTALL
|
|||
include LICENSE
|
||||
include MANIFEST.in
|
||||
include scrapy/mime.types
|
||||
include scrapyd/default_scrapyd.conf
|
||||
recursive-include scrapyd/tests *.egg
|
||||
recursive-include scrapy/templates *
|
||||
recursive-include scrapy/tests/sample_data *
|
||||
recursive-include scrapy license.txt
|
||||
|
|
|
|||
|
|
@ -1,8 +1,18 @@
|
|||
TRIAL := $(shell which trial)
|
||||
BRANCH := $(shell git rev-parse --abbrev-ref HEAD)
|
||||
ifeq ($(BRANCH),master)
|
||||
export SCRAPY_VERSION_FROM_GIT=1
|
||||
endif
|
||||
export PYTHONPATH=$(PWD)
|
||||
|
||||
deb:
|
||||
test:
|
||||
coverage run --branch $(TRIAL) --reporter=text scrapy.tests
|
||||
rm -rf htmlcov && coverage html
|
||||
-s3cmd sync -P htmlcov/ s3://static.scrapy.org/coverage-scrapy-$(BRANCH)/
|
||||
|
||||
build:
|
||||
python extras/makedeb.py build
|
||||
|
||||
deb-clean:
|
||||
debclean
|
||||
python extras/makedeb.py clean
|
||||
clean:
|
||||
git checkout debian scrapy/__init__.py
|
||||
git clean -dfq
|
||||
|
|
|
|||
|
|
@ -37,7 +37,7 @@ fi
|
|||
|
||||
find . -name '*.py[co]' -delete
|
||||
if [ $# -eq 0 ]; then
|
||||
$trial --reporter=text scrapy scrapyd
|
||||
$trial --reporter=text scrapy
|
||||
else
|
||||
$trial "$@"
|
||||
fi
|
||||
|
|
|
|||
14
bin/scrapyd
14
bin/scrapyd
|
|
@ -1,14 +0,0 @@
|
|||
#!/bin/sh
|
||||
|
||||
repotac=$(cd $(dirname $0)/../extras; pwd)/scrapyd.tac
|
||||
|
||||
if [ -f "$repotac" ]; then
|
||||
tacfile="$repotac"
|
||||
elif [ -f "/usr/share/scrapyd/scrapyd.tac" ]; then
|
||||
tacfile="/usr/share/scrapyd/scrapyd.tac"
|
||||
else
|
||||
echo "Unable to find scrapy.tac file"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
twistd -ny "$tacfile"
|
||||
|
|
@ -17,14 +17,3 @@ Description: Python web crawling and scraping framework
|
|||
used to crawl websites and extract structured data from their pages.
|
||||
It can be used for a wide range of purposes, from data mining to
|
||||
monitoring and automated testing.
|
||||
|
||||
Package: scrapyd-SUFFIX
|
||||
Architecture: all
|
||||
Depends: scrapy-SUFFIX (>= ${source:Version}), python-setuptools
|
||||
Conflicts: scrapyd, scrapyd-0.11
|
||||
Provides: scrapyd
|
||||
Description: Scrapy Service
|
||||
The Scrapy service allows you to deploy your Scrapy projects by building
|
||||
Python eggs of them and uploading them to the Scrapy service using a JSON API
|
||||
that you can also use for scheduling spider runs. It supports multiple
|
||||
projects also.
|
||||
|
|
|
|||
|
|
@ -3,6 +3,3 @@
|
|||
|
||||
%:
|
||||
dh $@
|
||||
|
||||
override_dh_installinit:
|
||||
dh_installinit --name scrapyd
|
||||
|
|
|
|||
|
|
@ -1,3 +1 @@
|
|||
usr/lib/python*/*-packages/scrapy*
|
||||
usr/bin
|
||||
extras/scrapy_bash_completion etc/bash_completion.d/
|
||||
|
|
|
|||
|
|
@ -1,8 +0,0 @@
|
|||
[scrapyd]
|
||||
http_port = 6800
|
||||
debug = off
|
||||
#max_proc = 1
|
||||
eggs_dir = /var/lib/scrapyd/eggs
|
||||
dbs_dir = /var/lib/scrapyd/dbs
|
||||
items_dir = /var/lib/scrapyd/items
|
||||
logs_dir = /var/log/scrapyd
|
||||
|
|
@ -1 +0,0 @@
|
|||
# Defaults for Scrapy service
|
||||
|
|
@ -1,5 +0,0 @@
|
|||
var/lib/scrapyd
|
||||
var/lib/scrapyd/eggs
|
||||
var/lib/scrapyd/dbs
|
||||
var/lib/scrapyd/items
|
||||
var/log/scrapyd
|
||||
|
|
@ -1,2 +0,0 @@
|
|||
debian/scrapyd-files/000-default etc/scrapyd/conf.d
|
||||
extras/scrapyd.tac usr/share/scrapyd
|
||||
|
|
@ -1,2 +0,0 @@
|
|||
new-package-should-close-itp-bug
|
||||
script-in-etc-init.d-not-registered-via-update-rc.d /etc/init.d/scrapyd
|
||||
|
|
@ -1,34 +0,0 @@
|
|||
#!/bin/sh
|
||||
|
||||
set -e
|
||||
|
||||
case "$1" in
|
||||
configure)
|
||||
# Create user to run the service as
|
||||
if [ -z "`id -u scrapy 2> /dev/null`" ]; then
|
||||
adduser --system --home /var/lib/scrapyd --gecos "scrapy" \
|
||||
--no-create-home --disabled-password \
|
||||
--quiet scrapy || true
|
||||
fi
|
||||
|
||||
chown scrapy:nogroup /var/log/scrapyd \
|
||||
/var/lib/scrapyd /var/lib/scrapyd/eggs /var/lib/scrapyd/dbs \
|
||||
/var/lib/scrapyd/items
|
||||
|
||||
update-python-modules -p # so upstart restart uses the new code
|
||||
;;
|
||||
|
||||
abort-upgrade|abort-remove|abort-deconfigure)
|
||||
;;
|
||||
|
||||
*)
|
||||
echo "postinst called with unknown argument \`$1'" >&2
|
||||
exit 1
|
||||
;;
|
||||
esac
|
||||
|
||||
#DEBHELPER#
|
||||
|
||||
exit 0
|
||||
|
||||
|
||||
|
|
@ -1,17 +0,0 @@
|
|||
#!/bin/sh
|
||||
|
||||
set -e
|
||||
|
||||
if [ purge = "$1" ]; then
|
||||
if [ -x "$(command -v deluser)" ]; then
|
||||
deluser --quiet --system scrapy > /dev/null || true
|
||||
else
|
||||
echo >&2 "not removing scrapy system account because deluser command was not found"
|
||||
fi
|
||||
fi
|
||||
|
||||
#DEBHELPER#
|
||||
|
||||
exit 0
|
||||
|
||||
|
||||
|
|
@ -1,13 +0,0 @@
|
|||
# Scrapy service
|
||||
|
||||
start on runlevel [2345]
|
||||
stop on runlevel [06]
|
||||
|
||||
script
|
||||
[ -r /etc/default/scrapyd ] && . /etc/default/scrapyd
|
||||
logdir=/var/log/scrapyd
|
||||
exec twistd -ny /usr/share/scrapyd/scrapyd.tac \
|
||||
-u scrapy -g nogroup \
|
||||
--pidfile /var/run/scrapyd.pid \
|
||||
-l $logdir/scrapyd.log >$logdir/scrapyd.out 2>$logdir/scrapyd.err
|
||||
end script
|
||||
|
|
@ -249,24 +249,6 @@ Usage examples::
|
|||
[FAILED] first_spider:parse
|
||||
>>> Returned 92 requests, expected 0..4
|
||||
|
||||
.. command:: server
|
||||
|
||||
server
|
||||
------
|
||||
|
||||
* Syntax: ``scrapy server``
|
||||
* Requires project: *yes*
|
||||
|
||||
Start Scrapyd server for this project, which can be referred from the JSON API
|
||||
with the project name ``default``. For more info see: :ref:`topics-scrapyd`.
|
||||
|
||||
Usage example::
|
||||
|
||||
$ scrapy server
|
||||
[ ... scrapyd starts and stays idle waiting for spiders to get scheduled ... ]
|
||||
|
||||
To schedule spiders, use the Scrapyd JSON API.
|
||||
|
||||
.. command:: list
|
||||
|
||||
list
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@ import sys, os, glob, shutil
|
|||
from subprocess import check_call
|
||||
|
||||
def build(suffix):
|
||||
for ifn in glob.glob("debian/scrapy.*") + glob.glob("debian/scrapyd.*"):
|
||||
for ifn in glob.glob("debian/scrapy.*"):
|
||||
s = open(ifn).read()
|
||||
s = s.replace('SUFFIX', suffix)
|
||||
pre, suf = ifn.split('.', 1)
|
||||
|
|
@ -21,8 +21,7 @@ def build(suffix):
|
|||
check_call('debuild -us -uc -b', shell=True)
|
||||
|
||||
def clean(suffix):
|
||||
for f in glob.glob("debian/python-scrapy%s*" % suffix) + \
|
||||
glob.glob("debian/scrapyd%s*" % suffix):
|
||||
for f in glob.glob("debian/python-scrapy%s*" % suffix):
|
||||
if os.path.isdir(f):
|
||||
shutil.rmtree(f)
|
||||
else:
|
||||
|
|
|
|||
|
|
@ -1,2 +0,0 @@
|
|||
from scrapyd import get_application
|
||||
application = get_application()
|
||||
|
|
@ -1,122 +0,0 @@
|
|||
#!/bin/bash
|
||||
#
|
||||
# This script is a quick system test for Scrapyd that:
|
||||
#
|
||||
# 1. runs scrapyd
|
||||
# 2. creates a new project and deploys it on scrapyd
|
||||
# 3. schedules a spider on scrapyd and waits for it to finish
|
||||
# 4. check the spider scraped the expected data
|
||||
#
|
||||
|
||||
set -e
|
||||
|
||||
export PATH=$PATH:$(pwd)/bin
|
||||
export PYTHONPATH=$PYTHONPATH:$(pwd)
|
||||
|
||||
scrapyd_dir=$(mktemp /tmp/test-scrapyd.XXXXXXX -d)
|
||||
scrapyd_log=$scrapyd_dir/scrapyd.log
|
||||
scrapy_dir=$(mktemp /tmp/test-scrapy.XXXXXXX -d)
|
||||
|
||||
echo "scrapyd dir: $scrapyd_dir"
|
||||
echo "scrapy dir : $scrapy_dir"
|
||||
|
||||
twistd -ny extras/scrapyd.tac -d $scrapyd_dir -l $scrapyd_log &
|
||||
|
||||
cd $scrapy_dir
|
||||
scrapy startproject testproj
|
||||
cd testproj
|
||||
cat > testproj/spiders/insophia.py <<!
|
||||
from scrapy.spider import BaseSpider
|
||||
from scrapy.selector import HtmlXPathSelector
|
||||
from scrapy.http import Request
|
||||
from scrapy.item import Item, Field
|
||||
|
||||
|
||||
class Section(Item):
|
||||
url = Field()
|
||||
title = Field()
|
||||
size = Field()
|
||||
arg = Field()
|
||||
|
||||
|
||||
class InsophiaSpider(BaseSpider):
|
||||
name = 'insophia'
|
||||
start_urls = ['http://insophia.com']
|
||||
|
||||
def __init__(self, *a, **kw):
|
||||
self.arg = kw.pop('arg')
|
||||
super(InsophiaSpider, self).__init__(*a, **kw)
|
||||
|
||||
def parse(self, response):
|
||||
hxs = HtmlXPathSelector(response)
|
||||
for url in hxs.select('//ul[@id="navlist"]/li/a/@href').extract():
|
||||
yield Request(url, callback=self.parse_section)
|
||||
|
||||
def parse_section(self, response):
|
||||
hxs = HtmlXPathSelector(response)
|
||||
title = hxs.select("//h2/text()").extract()[0]
|
||||
yield Section(url=response.url, title=title, size=len(response.body),
|
||||
arg=self.arg)
|
||||
!
|
||||
|
||||
cat > scrapy.cfg <<!
|
||||
[settings]
|
||||
default = testproj.settings
|
||||
|
||||
[deploy]
|
||||
url = http://localhost:6800/
|
||||
project = testproj
|
||||
!
|
||||
|
||||
scrapy deploy
|
||||
|
||||
curl -s http://localhost:6800/schedule.json -d project=testproj -d spider=insophia -d arg=SOME_ARGUMENT
|
||||
|
||||
echo "waiting 20 seconds for spider to run and finish..."
|
||||
sleep 20
|
||||
|
||||
kill %1
|
||||
wait %1
|
||||
|
||||
if ! grep -q "Process finished" $scrapyd_log; then
|
||||
echo "error: 'Process finished' not found on scrapyd log"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
feed_path=$(find $scrapyd_dir/items -name '*.jl')
|
||||
if [ ! -f "$feed_path" ]; then
|
||||
echo "items feed not generated: $feed_path"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
log_path=$(find $scrapyd_dir/logs -name '*.log')
|
||||
if [ ! -f "$log_path" ]; then
|
||||
echo "log file not generated: $log_path"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
numitems="$(cat $feed_path | wc -l)"
|
||||
if [ "$numitems" != "6" ]; then
|
||||
echo "error: wrong number of items scraped: $numitems"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
numscraped="$(cat $log_path | grep Scraped | wc -l)"
|
||||
if [ "$numscraped" != "6" ]; then
|
||||
echo "error: wrong number of 'Scraped' lines in log: $numscraped"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
if ! grep -q "About Us" $feed_path; then
|
||||
echo "error: About Us page not scraped"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
if ! grep -q "SOME_ARGUMENT" $feed_path; then
|
||||
echo "error: spider argument not found in scraped items"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
rm -rf /tmp/test-scrapyd.* /tmp/test-scrapy.*
|
||||
|
||||
echo "All tests OK"
|
||||
|
|
@ -1,22 +0,0 @@
|
|||
from __future__ import absolute_import
|
||||
|
||||
from scrapy.command import ScrapyCommand
|
||||
from scrapy.exceptions import UsageError
|
||||
|
||||
class Command(ScrapyCommand):
|
||||
|
||||
requires_project = True
|
||||
|
||||
def short_desc(self):
|
||||
return "Start Scrapyd server for this project"
|
||||
|
||||
def long_desc(self):
|
||||
return "Start Scrapyd server for this project, which can be referred " \
|
||||
"from the JSON API with the name 'default'"
|
||||
|
||||
def run(self, args, opts):
|
||||
try:
|
||||
from scrapyd.script import execute
|
||||
execute()
|
||||
except ImportError:
|
||||
raise UsageError("Scrapyd is not available in this system")
|
||||
|
|
@ -1,9 +0,0 @@
|
|||
from scrapy.utils.misc import load_object
|
||||
from .config import Config
|
||||
|
||||
def get_application(config=None):
|
||||
if config is None:
|
||||
config = Config()
|
||||
apppath = config.get('application', 'scrapyd.app.application')
|
||||
appfunc = load_object(apppath)
|
||||
return appfunc(config)
|
||||
|
|
@ -1,44 +0,0 @@
|
|||
from twisted.application.service import Application
|
||||
from twisted.application.internet import TimerService, TCPServer
|
||||
from twisted.web import server
|
||||
from twisted.python import log
|
||||
|
||||
from scrapy.utils.misc import load_object
|
||||
|
||||
from .interfaces import IEggStorage, IPoller, ISpiderScheduler, IEnvironment
|
||||
from .eggstorage import FilesystemEggStorage
|
||||
from .scheduler import SpiderScheduler
|
||||
from .poller import QueuePoller
|
||||
from .environ import Environment
|
||||
from .website import Root
|
||||
from .config import Config
|
||||
|
||||
def application(config):
|
||||
app = Application("Scrapyd")
|
||||
http_port = config.getint('http_port', 6800)
|
||||
bind_address = config.get('bind_address', '0.0.0.0')
|
||||
|
||||
poller = QueuePoller(config)
|
||||
eggstorage = FilesystemEggStorage(config)
|
||||
scheduler = SpiderScheduler(config)
|
||||
environment = Environment(config)
|
||||
|
||||
app.setComponent(IPoller, poller)
|
||||
app.setComponent(IEggStorage, eggstorage)
|
||||
app.setComponent(ISpiderScheduler, scheduler)
|
||||
app.setComponent(IEnvironment, environment)
|
||||
|
||||
laupath = config.get('launcher', 'scrapyd.launcher.Launcher')
|
||||
laucls = load_object(laupath)
|
||||
launcher = laucls(config, app)
|
||||
|
||||
timer = TimerService(5, poller.poll)
|
||||
webservice = TCPServer(http_port, server.Site(Root(config, app)), interface=bind_address)
|
||||
log.msg(format="Scrapyd web console available at http://%(bind_address)s:%(http_port)s/",
|
||||
bind_address=bind_address, http_port=http_port)
|
||||
|
||||
launcher.setServiceParent(app)
|
||||
timer.setServiceParent(app)
|
||||
webservice.setServiceParent(app)
|
||||
|
||||
return app
|
||||
|
|
@ -1,62 +0,0 @@
|
|||
import glob
|
||||
from cStringIO import StringIO
|
||||
from pkgutil import get_data
|
||||
from ConfigParser import SafeConfigParser, NoSectionError, NoOptionError
|
||||
|
||||
from scrapy.utils.conf import closest_scrapy_cfg
|
||||
|
||||
class Config(object):
|
||||
"""A ConfigParser wrapper to support defaults when calling instance
|
||||
methods, and also tied to a single section"""
|
||||
|
||||
SECTION = 'scrapyd'
|
||||
|
||||
def __init__(self, values=None, extra_sources=()):
|
||||
if values is None:
|
||||
sources = self._getsources()
|
||||
default_config = get_data(__package__, 'default_scrapyd.conf')
|
||||
self.cp = SafeConfigParser()
|
||||
self.cp.readfp(StringIO(default_config))
|
||||
self.cp.read(sources)
|
||||
for fp in extra_sources:
|
||||
self.cp.readfp(fp)
|
||||
else:
|
||||
self.cp = SafeConfigParser(values)
|
||||
self.cp.add_section(self.SECTION)
|
||||
|
||||
def _getsources(self):
|
||||
sources = ['/etc/scrapyd/scrapyd.conf', r'c:\scrapyd\scrapyd.conf']
|
||||
sources += sorted(glob.glob('/etc/scrapyd/conf.d/*'))
|
||||
sources += ['scrapyd.conf']
|
||||
scrapy_cfg = closest_scrapy_cfg()
|
||||
if scrapy_cfg:
|
||||
sources.append(scrapy_cfg)
|
||||
return sources
|
||||
|
||||
def _getany(self, method, option, default):
|
||||
try:
|
||||
return method(self.SECTION, option)
|
||||
except (NoSectionError, NoOptionError):
|
||||
if default is not None:
|
||||
return default
|
||||
raise
|
||||
|
||||
def get(self, option, default=None):
|
||||
return self._getany(self.cp.get, option, default)
|
||||
|
||||
def getint(self, option, default=None):
|
||||
return self._getany(self.cp.getint, option, default)
|
||||
|
||||
def getfloat(self, option, default=None):
|
||||
return self._getany(self.cp.getfloat, option, default)
|
||||
|
||||
def getboolean(self, option, default=None):
|
||||
return self._getany(self.cp.getboolean, option, default)
|
||||
|
||||
def items(self, section, default=None):
|
||||
try:
|
||||
return self.cp.items(section)
|
||||
except (NoSectionError, NoOptionError):
|
||||
if default is not None:
|
||||
return default
|
||||
raise
|
||||
|
|
@ -1,25 +0,0 @@
|
|||
[scrapyd]
|
||||
eggs_dir = eggs
|
||||
logs_dir = logs
|
||||
items_dir = items
|
||||
jobs_to_keep = 5
|
||||
dbs_dir = dbs
|
||||
max_proc = 0
|
||||
max_proc_per_cpu = 4
|
||||
finished_to_keep = 100
|
||||
http_port = 6800
|
||||
debug = off
|
||||
runner = scrapyd.runner
|
||||
application = scrapyd.app.application
|
||||
launcher = scrapyd.launcher.Launcher
|
||||
|
||||
[services]
|
||||
schedule.json = scrapyd.webservice.Schedule
|
||||
cancel.json = scrapyd.webservice.Cancel
|
||||
addversion.json = scrapyd.webservice.AddVersion
|
||||
listprojects.json = scrapyd.webservice.ListProjects
|
||||
listversions.json = scrapyd.webservice.ListVersions
|
||||
listspiders.json = scrapyd.webservice.ListSpiders
|
||||
delproject.json = scrapyd.webservice.DeleteProject
|
||||
delversion.json = scrapyd.webservice.DeleteVersion
|
||||
listjobs.json = scrapyd.webservice.ListJobs
|
||||
|
|
@ -1,49 +0,0 @@
|
|||
from glob import glob
|
||||
from os import path, makedirs, remove
|
||||
from shutil import copyfileobj, rmtree
|
||||
from distutils.version import LooseVersion
|
||||
|
||||
from zope.interface import implements
|
||||
|
||||
from .interfaces import IEggStorage
|
||||
|
||||
class FilesystemEggStorage(object):
|
||||
|
||||
implements(IEggStorage)
|
||||
|
||||
def __init__(self, config):
|
||||
self.basedir = config.get('eggs_dir', 'eggs')
|
||||
|
||||
def put(self, eggfile, project, version):
|
||||
eggpath = self._eggpath(project, version)
|
||||
eggdir = path.dirname(eggpath)
|
||||
if not path.exists(eggdir):
|
||||
makedirs(eggdir)
|
||||
with open(eggpath, 'wb') as f:
|
||||
copyfileobj(eggfile, f)
|
||||
|
||||
def get(self, project, version=None):
|
||||
if version is None:
|
||||
try:
|
||||
version = self.list(project)[-1]
|
||||
except IndexError:
|
||||
return None, None
|
||||
return version, open(self._eggpath(project, version), 'rb')
|
||||
|
||||
def list(self, project):
|
||||
eggdir = path.join(self.basedir, project)
|
||||
versions = [path.splitext(path.basename(x))[0] \
|
||||
for x in glob("%s/*.egg" % eggdir)]
|
||||
return sorted(versions, key=LooseVersion)
|
||||
|
||||
def delete(self, project, version=None):
|
||||
if version is None:
|
||||
rmtree(path.join(self.basedir, project))
|
||||
else:
|
||||
remove(self._eggpath(project, version))
|
||||
if not self.list(project): # remove project if no versions left
|
||||
self.delete(project)
|
||||
|
||||
def _eggpath(self, project, version):
|
||||
x = path.join(self.basedir, project, "%s.egg" % version)
|
||||
return x
|
||||
|
|
@ -1,14 +0,0 @@
|
|||
import os, pkg_resources
|
||||
|
||||
def activate_egg(eggpath):
|
||||
"""Activate a Scrapy egg file. This is meant to be used from egg runners
|
||||
to activate a Scrapy egg file. Don't use it from other code as it may
|
||||
leave unwanted side effects.
|
||||
"""
|
||||
try:
|
||||
d = pkg_resources.find_distributions(eggpath).next()
|
||||
except StopIteration:
|
||||
raise ValueError("Unknown or corrupt egg")
|
||||
d.activate()
|
||||
settings_module = d.get_entry_info('scrapy', 'settings').module_name
|
||||
os.environ.setdefault('SCRAPY_SETTINGS_MODULE', settings_module)
|
||||
|
|
@ -1,46 +0,0 @@
|
|||
import os
|
||||
|
||||
from zope.interface import implements
|
||||
|
||||
from .interfaces import IEnvironment
|
||||
|
||||
class Environment(object):
|
||||
|
||||
implements(IEnvironment)
|
||||
|
||||
def __init__(self, config, initenv=os.environ):
|
||||
self.dbs_dir = config.get('dbs_dir', 'dbs')
|
||||
self.logs_dir = config.get('logs_dir', 'logs')
|
||||
self.items_dir = config.get('items_dir', 'items')
|
||||
self.jobs_to_keep = config.getint('jobs_to_keep', 5)
|
||||
if config.cp.has_section('settings'):
|
||||
self.settings = dict(config.cp.items('settings'))
|
||||
else:
|
||||
self.settings = {}
|
||||
self.initenv = initenv
|
||||
|
||||
def get_environment(self, message, slot):
|
||||
project = message['_project']
|
||||
env = self.initenv.copy()
|
||||
env['SCRAPY_SLOT'] = str(slot)
|
||||
env['SCRAPY_PROJECT'] = project
|
||||
env['SCRAPY_SPIDER'] = message['_spider']
|
||||
env['SCRAPY_JOB'] = message['_job']
|
||||
if project in self.settings:
|
||||
env['SCRAPY_SETTINGS_MODULE'] = self.settings[project]
|
||||
if self.logs_dir:
|
||||
env['SCRAPY_LOG_FILE'] = self._get_file(message, self.logs_dir, 'log')
|
||||
if self.items_dir:
|
||||
env['SCRAPY_FEED_URI'] = self._get_file(message, self.items_dir, 'jl')
|
||||
return env
|
||||
|
||||
def _get_file(self, message, dir, ext):
|
||||
logsdir = os.path.join(dir, message['_project'], \
|
||||
message['_spider'])
|
||||
if not os.path.exists(logsdir):
|
||||
os.makedirs(logsdir)
|
||||
to_delete = sorted((os.path.join(logsdir, x) for x in \
|
||||
os.listdir(logsdir)), key=os.path.getmtime)[:-self.jobs_to_keep]
|
||||
for x in to_delete:
|
||||
os.remove(x)
|
||||
return os.path.join(logsdir, "%s.%s" % (message['_job'], ext))
|
||||
|
|
@ -1,109 +0,0 @@
|
|||
from zope.interface import Interface
|
||||
|
||||
class IEggStorage(Interface):
|
||||
"""A component that handles storing and retrieving eggs"""
|
||||
|
||||
def put(eggfile, project, version):
|
||||
"""Store the egg (passed in the file object) under the given project and
|
||||
version"""
|
||||
|
||||
def get(project, version=None):
|
||||
"""Return a tuple (version, file) with the the egg for the specified
|
||||
project and version. If version is None, the latest version is
|
||||
returned. If no egg is found for the given project/version (None, None)
|
||||
should be returned."""
|
||||
|
||||
def list(project):
|
||||
"""Return the list of versions which have eggs stored (for the given
|
||||
project) in order (the latest version is the currently used)."""
|
||||
|
||||
def delete(project, version=None):
|
||||
"""Delete the egg stored for the given project and version. If should
|
||||
also delete the project if no versions are left"""
|
||||
|
||||
|
||||
class IPoller(Interface):
|
||||
"""A component that polls for projects that need to run"""
|
||||
|
||||
def poll():
|
||||
"""Called periodically to poll for projects"""
|
||||
|
||||
def next():
|
||||
"""Return the next message.
|
||||
|
||||
It should return a Deferred which will get fired when there is a new
|
||||
project that needs to run, or already fired if there was a project
|
||||
waiting to run already.
|
||||
|
||||
The message is a dict containing (at least):
|
||||
* the name of the project to be run in the '_project' key
|
||||
* the name of the spider to be run in the '_spider' key
|
||||
* a unique identifier for this run in the `_job` key
|
||||
This message will be passed later to IEnvironment.get_environment().
|
||||
"""
|
||||
|
||||
def update_projects():
|
||||
"""Called when projects may have changed, to refresh the available
|
||||
projects"""
|
||||
|
||||
|
||||
class ISpiderQueue(Interface):
|
||||
|
||||
def add(name, **spider_args):
|
||||
"""Add a spider to the queue given its name a some spider arguments.
|
||||
|
||||
This method can return a deferred. """
|
||||
|
||||
def pop():
|
||||
"""Pop the next mesasge from the queue. The messages is a dict
|
||||
conaining a key 'name' with the spider name and other keys as spider
|
||||
attributes.
|
||||
|
||||
This method can return a deferred. """
|
||||
|
||||
def list():
|
||||
"""Return a list with the messages in the queue. Each message is a dict
|
||||
which must have a 'name' key (with the spider name), and other optional
|
||||
keys that will be used as spider arguments, to create the spider.
|
||||
|
||||
This method can return a deferred. """
|
||||
|
||||
def count():
|
||||
"""Return the number of spiders in the queue.
|
||||
|
||||
This method can return a deferred. """
|
||||
|
||||
def remove(func):
|
||||
"""Remove all elements from the queue for which func(element) is true,
|
||||
and return the number of removed elements.
|
||||
"""
|
||||
|
||||
def clear():
|
||||
"""Clear the queue.
|
||||
|
||||
This method can return a deferred. """
|
||||
|
||||
|
||||
class ISpiderScheduler(Interface):
|
||||
"""A component to schedule spider runs"""
|
||||
|
||||
def schedule(project, spider_name, **spider_args):
|
||||
"""Schedule a spider for the given project"""
|
||||
|
||||
def list_projects():
|
||||
"""Return the list of available projects"""
|
||||
|
||||
def update_projects():
|
||||
"""Called when projects may have changed, to refresh the available
|
||||
projects"""
|
||||
|
||||
|
||||
class IEnvironment(Interface):
|
||||
"""A component to generate the environment of crawler processes"""
|
||||
|
||||
def get_environment(message, slot):
|
||||
"""Return the environment variables to use for running the process.
|
||||
|
||||
`message` is the message received from the IPoller.next() method
|
||||
`slot` is the Launcher slot where the process will be running.
|
||||
"""
|
||||
|
|
@ -1,102 +0,0 @@
|
|||
import sys
|
||||
from datetime import datetime
|
||||
from multiprocessing import cpu_count
|
||||
|
||||
from twisted.internet import reactor, defer, protocol, error
|
||||
from twisted.application.service import Service
|
||||
from twisted.python import log
|
||||
|
||||
from scrapy.utils.python import stringify_dict
|
||||
from scrapyd.utils import get_crawl_args
|
||||
from .interfaces import IPoller, IEnvironment
|
||||
|
||||
class Launcher(Service):
|
||||
|
||||
name = 'launcher'
|
||||
|
||||
def __init__(self, config, app):
|
||||
self.processes = {}
|
||||
self.finished = []
|
||||
self.finished_to_keep = config.getint('finished_to_keep', 100)
|
||||
self.max_proc = self._get_max_proc(config)
|
||||
self.runner = config.get('runner', 'scrapyd.runner')
|
||||
self.app = app
|
||||
|
||||
def startService(self):
|
||||
for slot in range(self.max_proc):
|
||||
self._wait_for_project(slot)
|
||||
log.msg(format='%(parent)s started: max_proc=%(max_proc)r, runner=%(runner)r',
|
||||
parent=self.parent.name, max_proc=self.max_proc,
|
||||
runner=self.runner, system='Launcher')
|
||||
|
||||
def _wait_for_project(self, slot):
|
||||
poller = self.app.getComponent(IPoller)
|
||||
poller.next().addCallback(self._spawn_process, slot)
|
||||
|
||||
def _spawn_process(self, message, slot):
|
||||
msg = stringify_dict(message, keys_only=False)
|
||||
project = msg['_project']
|
||||
args = [sys.executable, '-m', self.runner, 'crawl']
|
||||
args += get_crawl_args(msg)
|
||||
e = self.app.getComponent(IEnvironment)
|
||||
env = e.get_environment(msg, slot)
|
||||
env = stringify_dict(env, keys_only=False)
|
||||
pp = ScrapyProcessProtocol(slot, project, msg['_spider'], \
|
||||
msg['_job'], env)
|
||||
pp.deferred.addBoth(self._process_finished, slot)
|
||||
reactor.spawnProcess(pp, sys.executable, args=args, env=env)
|
||||
self.processes[slot] = pp
|
||||
|
||||
def _process_finished(self, _, slot):
|
||||
process = self.processes.pop(slot)
|
||||
process.end_time = datetime.now()
|
||||
self.finished.append(process)
|
||||
del self.finished[:-self.finished_to_keep] # keep last 100 finished jobs
|
||||
self._wait_for_project(slot)
|
||||
|
||||
def _get_max_proc(self, config):
|
||||
max_proc = config.getint('max_proc', 0)
|
||||
if not max_proc:
|
||||
try:
|
||||
cpus = cpu_count()
|
||||
except NotImplementedError:
|
||||
cpus = 1
|
||||
max_proc = cpus * config.getint('max_proc_per_cpu', 4)
|
||||
return max_proc
|
||||
|
||||
class ScrapyProcessProtocol(protocol.ProcessProtocol):
|
||||
|
||||
def __init__(self, slot, project, spider, job, env):
|
||||
self.slot = slot
|
||||
self.pid = None
|
||||
self.project = project
|
||||
self.spider = spider
|
||||
self.job = job
|
||||
self.start_time = datetime.now()
|
||||
self.end_time = None
|
||||
self.env = env
|
||||
self.logfile = env.get('SCRAPY_LOG_FILE')
|
||||
self.itemsfile = env.get('SCRAPY_FEED_URI')
|
||||
self.deferred = defer.Deferred()
|
||||
|
||||
def outReceived(self, data):
|
||||
log.msg(data.rstrip(), system="Launcher,%d/stdout" % self.pid)
|
||||
|
||||
def errReceived(self, data):
|
||||
log.msg(data.rstrip(), system="Launcher,%d/stderr" % self.pid)
|
||||
|
||||
def connectionMade(self):
|
||||
self.pid = self.transport.pid
|
||||
self.log("Process started: ")
|
||||
|
||||
def processEnded(self, status):
|
||||
if isinstance(status.value, error.ProcessDone):
|
||||
self.log("Process finished: ")
|
||||
else:
|
||||
self.log("Process died: exitstatus=%r " % status.value.exitCode)
|
||||
self.deferred.callback(self)
|
||||
|
||||
def log(self, action):
|
||||
fmt = '%(action)s project=%(project)r spider=%(spider)r job=%(job)r pid=%(pid)r log=%(log)r items=%(items)r'
|
||||
log.msg(format=fmt, action=action, project=self.project, spider=self.spider,
|
||||
job=self.job, pid=self.pid, log=self.logfile, items=self.itemsfile)
|
||||
|
|
@ -1,36 +0,0 @@
|
|||
from zope.interface import implements
|
||||
from twisted.internet.defer import DeferredQueue, inlineCallbacks, maybeDeferred, returnValue
|
||||
|
||||
from .utils import get_spider_queues
|
||||
from .interfaces import IPoller
|
||||
|
||||
class QueuePoller(object):
|
||||
|
||||
implements(IPoller)
|
||||
|
||||
def __init__(self, config):
|
||||
self.config = config
|
||||
self.update_projects()
|
||||
self.dq = DeferredQueue(size=1)
|
||||
|
||||
@inlineCallbacks
|
||||
def poll(self):
|
||||
if self.dq.pending:
|
||||
return
|
||||
for p, q in self.queues.iteritems():
|
||||
c = yield maybeDeferred(q.count)
|
||||
if c:
|
||||
msg = yield maybeDeferred(q.pop)
|
||||
returnValue(self.dq.put(self._message(msg, p)))
|
||||
|
||||
def next(self):
|
||||
return self.dq.get()
|
||||
|
||||
def update_projects(self):
|
||||
self.queues = get_spider_queues(self.config)
|
||||
|
||||
def _message(self, queue_msg, project):
|
||||
d = queue_msg.copy()
|
||||
d['_project'] = project
|
||||
d['_spider'] = d.pop('name')
|
||||
return d
|
||||
|
|
@ -1,39 +0,0 @@
|
|||
import sys
|
||||
import os
|
||||
import shutil
|
||||
import tempfile
|
||||
from contextlib import contextmanager
|
||||
|
||||
from scrapyd import get_application
|
||||
from scrapyd.interfaces import IEggStorage
|
||||
from scrapyd.eggutils import activate_egg
|
||||
|
||||
@contextmanager
|
||||
def project_environment(project):
|
||||
app = get_application()
|
||||
eggstorage = app.getComponent(IEggStorage)
|
||||
version, eggfile = eggstorage.get(project)
|
||||
if eggfile:
|
||||
prefix = '%s-%s-' % (project, version)
|
||||
fd, eggpath = tempfile.mkstemp(prefix=prefix, suffix='.egg')
|
||||
lf = os.fdopen(fd, 'wb')
|
||||
shutil.copyfileobj(eggfile, lf)
|
||||
lf.close()
|
||||
activate_egg(eggpath)
|
||||
else:
|
||||
eggpath = None
|
||||
try:
|
||||
assert 'scrapy.conf' not in sys.modules, "Scrapy settings already loaded"
|
||||
yield
|
||||
finally:
|
||||
if eggpath:
|
||||
os.remove(eggpath)
|
||||
|
||||
def main():
|
||||
project = os.environ['SCRAPY_PROJECT']
|
||||
with project_environment(project):
|
||||
from scrapy.cmdline import execute
|
||||
execute()
|
||||
|
||||
if __name__ == '__main__':
|
||||
main()
|
||||
|
|
@ -1,22 +0,0 @@
|
|||
from zope.interface import implements
|
||||
|
||||
from .interfaces import ISpiderScheduler
|
||||
from .utils import get_spider_queues
|
||||
|
||||
class SpiderScheduler(object):
|
||||
|
||||
implements(ISpiderScheduler)
|
||||
|
||||
def __init__(self, config):
|
||||
self.config = config
|
||||
self.update_projects()
|
||||
|
||||
def schedule(self, project, spider_name, **spider_args):
|
||||
q = self.queues[project]
|
||||
q.add(spider_name, **spider_args)
|
||||
|
||||
def list_projects(self):
|
||||
return self.queues.keys()
|
||||
|
||||
def update_projects(self):
|
||||
self.queues = get_spider_queues(self.config)
|
||||
|
|
@ -1,42 +0,0 @@
|
|||
"""This module can be used to execute Scrapyd from a Scrapy command"""
|
||||
|
||||
import sys
|
||||
import os
|
||||
from cStringIO import StringIO
|
||||
|
||||
from twisted.python import log
|
||||
from twisted.internet import reactor
|
||||
from twisted.application import app
|
||||
|
||||
from scrapy.utils.project import project_data_dir
|
||||
|
||||
from scrapyd import get_application
|
||||
from scrapyd.config import Config
|
||||
|
||||
def _get_config():
|
||||
datadir = os.path.join(project_data_dir(), 'scrapyd')
|
||||
conf = {
|
||||
'eggs_dir': os.path.join(datadir, 'eggs'),
|
||||
'logs_dir': os.path.join(datadir, 'logs'),
|
||||
'items_dir': os.path.join(datadir, 'items'),
|
||||
'dbs_dir': os.path.join(datadir, 'dbs'),
|
||||
}
|
||||
for k in ['eggs_dir', 'logs_dir', 'items_dir', 'dbs_dir']: # create dirs
|
||||
d = conf[k]
|
||||
if not os.path.exists(d):
|
||||
os.makedirs(d)
|
||||
scrapyd_conf = """
|
||||
[scrapyd]
|
||||
eggs_dir = %(eggs_dir)s
|
||||
logs_dir = %(logs_dir)s
|
||||
items_dir = %(items_dir)s
|
||||
dbs_dir = %(dbs_dir)s
|
||||
""" % conf
|
||||
return Config(extra_sources=[StringIO(scrapyd_conf)])
|
||||
|
||||
def execute():
|
||||
config = _get_config()
|
||||
log.startLogging(sys.stderr)
|
||||
application = get_application(config)
|
||||
app.startApplication(application, False)
|
||||
reactor.run()
|
||||
|
|
@ -1,33 +0,0 @@
|
|||
from zope.interface import implements
|
||||
|
||||
from scrapyd.interfaces import ISpiderQueue
|
||||
from scrapyd.sqlite import JsonSqlitePriorityQueue
|
||||
|
||||
|
||||
class SqliteSpiderQueue(object):
|
||||
|
||||
implements(ISpiderQueue)
|
||||
|
||||
def __init__(self, database=None, table='spider_queue'):
|
||||
self.q = JsonSqlitePriorityQueue(database, table)
|
||||
|
||||
def add(self, name, **spider_args):
|
||||
d = spider_args.copy()
|
||||
d['name'] = name
|
||||
priority = float(d.pop('priority', 0))
|
||||
self.q.put(d, priority)
|
||||
|
||||
def pop(self):
|
||||
return self.q.pop()
|
||||
|
||||
def count(self):
|
||||
return len(self.q)
|
||||
|
||||
def list(self):
|
||||
return [x[0] for x in self.q]
|
||||
|
||||
def remove(self, func):
|
||||
return self.q.remove(func)
|
||||
|
||||
def clear(self):
|
||||
self.q.clear()
|
||||
|
|
@ -1,172 +0,0 @@
|
|||
import sqlite3
|
||||
import cPickle
|
||||
import json
|
||||
from UserDict import DictMixin
|
||||
|
||||
|
||||
class SqliteDict(DictMixin):
|
||||
"""SQLite-backed dictionary"""
|
||||
|
||||
def __init__(self, database=None, table="dict"):
|
||||
self.database = database or ':memory:'
|
||||
self.table = table
|
||||
# about check_same_thread: http://twistedmatrix.com/trac/ticket/4040
|
||||
self.conn = sqlite3.connect(self.database, check_same_thread=False)
|
||||
q = "create table if not exists %s (key text primary key, value blob)" \
|
||||
% table
|
||||
self.conn.execute(q)
|
||||
|
||||
def __getitem__(self, key):
|
||||
key = self.encode(key)
|
||||
q = "select value from %s where key=?" % self.table
|
||||
value = self.conn.execute(q, (key,)).fetchone()
|
||||
if value:
|
||||
return self.decode(value[0])
|
||||
raise KeyError(key)
|
||||
|
||||
def __setitem__(self, key, value):
|
||||
key, value = self.encode(key), self.encode(value)
|
||||
q = "insert or replace into %s (key, value) values (?,?)" % self.table
|
||||
self.conn.execute(q, (key, value))
|
||||
self.conn.commit()
|
||||
|
||||
def __delitem__(self, key):
|
||||
key = self.encode(key)
|
||||
q = "delete from %s where key=?" % self.table
|
||||
self.conn.execute(q, (key,))
|
||||
self.conn.commit()
|
||||
|
||||
def iterkeys(self):
|
||||
q = "select key from %s" % self.table
|
||||
return (self.decode(x[0]) for x in self.conn.execute(q))
|
||||
|
||||
def keys(self):
|
||||
return list(self.iterkeys())
|
||||
|
||||
def itervalues(self):
|
||||
q = "select value from %s" % self.table
|
||||
return (self.decode(x[0]) for x in self.conn.execute(q))
|
||||
|
||||
def values(self):
|
||||
return list(self.itervalues())
|
||||
|
||||
def iteritems(self):
|
||||
q = "select key, value from %s" % self.table
|
||||
return ((self.decode(x[0]), self.decode(x[1])) for x in self.conn.execute(q))
|
||||
|
||||
def items(self):
|
||||
return list(self.iteritems())
|
||||
|
||||
def encode(self, obj):
|
||||
return obj
|
||||
|
||||
def decode(self, text):
|
||||
return text
|
||||
|
||||
|
||||
class PickleSqliteDict(SqliteDict):
|
||||
|
||||
def encode(self, obj):
|
||||
return buffer(cPickle.dumps(obj, protocol=2))
|
||||
|
||||
def decode(self, text):
|
||||
return cPickle.loads(str(text))
|
||||
|
||||
|
||||
class JsonSqliteDict(SqliteDict):
|
||||
|
||||
def encode(self, obj):
|
||||
return json.dumps(obj)
|
||||
|
||||
def decode(self, text):
|
||||
return json.loads(text)
|
||||
|
||||
|
||||
|
||||
class SqlitePriorityQueue(object):
|
||||
"""SQLite priority queue. It relies on SQLite concurrency support for
|
||||
providing atomic inter-process operations.
|
||||
"""
|
||||
|
||||
def __init__(self, database=None, table="queue"):
|
||||
self.database = database or ':memory:'
|
||||
self.table = table
|
||||
# about check_same_thread: http://twistedmatrix.com/trac/ticket/4040
|
||||
self.conn = sqlite3.connect(self.database, check_same_thread=False)
|
||||
q = "create table if not exists %s (id integer primary key, " \
|
||||
"priority real key, message blob)" % table
|
||||
self.conn.execute(q)
|
||||
|
||||
def put(self, message, priority=0.0):
|
||||
args = (priority, self.encode(message))
|
||||
q = "insert into %s (priority, message) values (?,?)" % self.table
|
||||
self.conn.execute(q, args)
|
||||
self.conn.commit()
|
||||
|
||||
def pop(self):
|
||||
q = "select id, message from %s order by priority desc limit 1" \
|
||||
% self.table
|
||||
idmsg = self.conn.execute(q).fetchone()
|
||||
if idmsg is None:
|
||||
return
|
||||
id, msg = idmsg
|
||||
q = "delete from %s where id=?" % self.table
|
||||
c = self.conn.execute(q, (id,))
|
||||
if not c.rowcount: # record vanished, so let's try again
|
||||
self.conn.rollback()
|
||||
return self.pop()
|
||||
self.conn.commit()
|
||||
return self.decode(msg)
|
||||
|
||||
def remove(self, func):
|
||||
q = "select id, message from %s" % self.table
|
||||
n = 0
|
||||
for id, msg in self.conn.execute(q):
|
||||
if func(self.decode(msg)):
|
||||
q = "delete from %s where id=?" % self.table
|
||||
c = self.conn.execute(q, (id,))
|
||||
if not c.rowcount: # record vanished, so let's try again
|
||||
self.conn.rollback()
|
||||
return self.remove(func)
|
||||
n += 1
|
||||
self.conn.commit()
|
||||
return n
|
||||
|
||||
def clear(self):
|
||||
self.conn.execute("delete from %s" % self.table)
|
||||
self.conn.commit()
|
||||
|
||||
def __len__(self):
|
||||
q = "select count(*) from %s" % self.table
|
||||
return self.conn.execute(q).fetchone()[0]
|
||||
|
||||
def __iter__(self):
|
||||
q = "select message, priority from %s order by priority desc" % \
|
||||
self.table
|
||||
return ((self.decode(x), y) for x, y in self.conn.execute(q))
|
||||
|
||||
def encode(self, obj):
|
||||
return obj
|
||||
|
||||
def decode(self, text):
|
||||
return text
|
||||
|
||||
|
||||
class PickleSqlitePriorityQueue(SqlitePriorityQueue):
|
||||
|
||||
def encode(self, obj):
|
||||
return buffer(cPickle.dumps(obj, protocol=2))
|
||||
|
||||
def decode(self, text):
|
||||
return cPickle.loads(str(text))
|
||||
|
||||
|
||||
class JsonSqlitePriorityQueue(SqlitePriorityQueue):
|
||||
|
||||
def encode(self, obj):
|
||||
return json.dumps(obj)
|
||||
|
||||
def decode(self, text):
|
||||
return json.loads(text)
|
||||
|
||||
|
||||
Binary file not shown.
|
|
@ -1,22 +0,0 @@
|
|||
import sys
|
||||
import unittest
|
||||
|
||||
class SettingsSafeModulesTest(unittest.TestCase):
|
||||
|
||||
# these modules must not load scrapy.conf
|
||||
SETTINGS_SAFE_MODULES = [
|
||||
'scrapy.utils.project',
|
||||
'scrapy.utils.conf',
|
||||
'scrapyd.interfaces',
|
||||
'scrapyd.eggutils',
|
||||
]
|
||||
|
||||
def test_modules_that_shouldnt_load_settings(self):
|
||||
sys.modules.pop('scrapy.conf', None)
|
||||
for m in self.SETTINGS_SAFE_MODULES:
|
||||
__import__(m)
|
||||
assert 'scrapy.conf' not in sys.modules, \
|
||||
"Module %r must not cause the scrapy.conf module to be loaded" % m
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
|
@ -1,43 +0,0 @@
|
|||
from cStringIO import StringIO
|
||||
|
||||
from twisted.trial import unittest
|
||||
|
||||
from zope.interface.verify import verifyObject
|
||||
|
||||
from scrapyd.interfaces import IEggStorage
|
||||
from scrapyd.config import Config
|
||||
from scrapyd.eggstorage import FilesystemEggStorage
|
||||
|
||||
class EggStorageTest(unittest.TestCase):
|
||||
|
||||
def setUp(self):
|
||||
d = self.mktemp()
|
||||
config = Config(values={'eggs_dir': d})
|
||||
self.eggst = FilesystemEggStorage(config)
|
||||
|
||||
def test_interface(self):
|
||||
verifyObject(IEggStorage, self.eggst)
|
||||
|
||||
def test_put_get_list_delete(self):
|
||||
self.eggst.put(StringIO("egg01"), 'mybot', '01')
|
||||
self.eggst.put(StringIO("egg03"), 'mybot', '03')
|
||||
self.eggst.put(StringIO("egg02"), 'mybot', '02')
|
||||
|
||||
self.assertEqual(self.eggst.list('mybot'), ['01', '02', '03'])
|
||||
self.assertEqual(self.eggst.list('mybot2'), [])
|
||||
|
||||
v, f = self.eggst.get('mybot')
|
||||
self.assertEqual(v, "03")
|
||||
self.assertEqual(f.read(), "egg03")
|
||||
f.close()
|
||||
|
||||
v, f = self.eggst.get('mybot', '02')
|
||||
self.assertEqual(v, "02")
|
||||
self.assertEqual(f.read(), "egg02")
|
||||
f.close()
|
||||
|
||||
self.eggst.delete('mybot', '02')
|
||||
self.assertEqual(self.eggst.list('mybot'), ['01', '03'])
|
||||
|
||||
self.eggst.delete('mybot')
|
||||
self.assertEqual(self.eggst.list('mybot'), [])
|
||||
|
|
@ -1,45 +0,0 @@
|
|||
import os
|
||||
|
||||
from twisted.trial import unittest
|
||||
|
||||
from zope.interface.verify import verifyObject
|
||||
|
||||
from scrapyd.interfaces import IEnvironment
|
||||
from scrapyd.config import Config
|
||||
from scrapyd.environ import Environment
|
||||
|
||||
class EnvironmentTest(unittest.TestCase):
|
||||
|
||||
def setUp(self):
|
||||
d = self.mktemp()
|
||||
os.mkdir(d)
|
||||
config = Config(values={'eggs_dir': d, 'logs_dir': d})
|
||||
config.cp.add_section('settings')
|
||||
config.cp.set('settings', 'newbot', 'newbot.settings')
|
||||
self.environ = Environment(config, initenv={})
|
||||
|
||||
def test_interface(self):
|
||||
verifyObject(IEnvironment, self.environ)
|
||||
|
||||
def test_get_environment_with_eggfile(self):
|
||||
msg = {'_project': 'mybot', '_spider': 'myspider', '_job': 'ID'}
|
||||
slot = 3
|
||||
env = self.environ.get_environment(msg, slot)
|
||||
self.assertEqual(env['SCRAPY_PROJECT'], 'mybot')
|
||||
self.assertEqual(env['SCRAPY_SLOT'], '3')
|
||||
self.assertEqual(env['SCRAPY_SPIDER'], 'myspider')
|
||||
self.assertEqual(env['SCRAPY_JOB'], 'ID')
|
||||
self.assert_(env['SCRAPY_LOG_FILE'].endswith(os.path.join('mybot', 'myspider', 'ID.log')))
|
||||
self.assert_(env['SCRAPY_FEED_URI'].endswith(os.path.join('mybot', 'myspider', 'ID.jl')))
|
||||
self.failIf('SCRAPY_SETTINGS_MODULE' in env)
|
||||
|
||||
def test_get_environment_with_no_items_dir(self):
|
||||
config = Config(values={'items_dir': '', 'logs_dir': ''})
|
||||
config.cp.add_section('settings')
|
||||
config.cp.set('settings', 'newbot', 'newbot.settings')
|
||||
msg = {'_project': 'mybot', '_spider': 'myspider', '_job': 'ID'}
|
||||
slot = 3
|
||||
environ = Environment(config, initenv={})
|
||||
env = environ.get_environment(msg, slot)
|
||||
self.failUnless('SCRAPY_FEED_URI' not in env)
|
||||
self.failUnless('SCRAPY_LOG_FILE' not in env)
|
||||
|
|
@ -1,41 +0,0 @@
|
|||
import os
|
||||
|
||||
from twisted.trial import unittest
|
||||
from twisted.internet.defer import Deferred
|
||||
|
||||
from zope.interface.verify import verifyObject
|
||||
|
||||
from scrapyd.interfaces import IPoller
|
||||
from scrapyd.config import Config
|
||||
from scrapyd.poller import QueuePoller
|
||||
from scrapyd.utils import get_spider_queues
|
||||
|
||||
class QueuePollerTest(unittest.TestCase):
|
||||
|
||||
def setUp(self):
|
||||
d = self.mktemp()
|
||||
eggs_dir = os.path.join(d, 'eggs')
|
||||
dbs_dir = os.path.join(d, 'dbs')
|
||||
os.makedirs(eggs_dir)
|
||||
os.makedirs(dbs_dir)
|
||||
os.makedirs(os.path.join(eggs_dir, 'mybot1'))
|
||||
os.makedirs(os.path.join(eggs_dir, 'mybot2'))
|
||||
config = Config(values={'eggs_dir': eggs_dir, 'dbs_dir': dbs_dir})
|
||||
self.queues = get_spider_queues(config)
|
||||
self.poller = QueuePoller(config)
|
||||
|
||||
def test_interface(self):
|
||||
verifyObject(IPoller, self.poller)
|
||||
|
||||
def test_poll_next(self):
|
||||
self.queues['mybot1'].add('spider1')
|
||||
self.queues['mybot2'].add('spider2')
|
||||
d1 = self.poller.next()
|
||||
d2 = self.poller.next()
|
||||
self.failUnless(isinstance(d1, Deferred))
|
||||
self.failIf(hasattr(d1, 'result'))
|
||||
self.poller.poll()
|
||||
self.queues['mybot1'].pop()
|
||||
self.poller.poll()
|
||||
self.failUnlessEqual(d1.result, {'_project': 'mybot1', '_spider': 'spider1'})
|
||||
self.failUnlessEqual(d2.result, {'_project': 'mybot2', '_spider': 'spider2'})
|
||||
|
|
@ -1,44 +0,0 @@
|
|||
import os
|
||||
|
||||
from twisted.trial import unittest
|
||||
|
||||
from zope.interface.verify import verifyObject
|
||||
|
||||
from scrapyd.interfaces import ISpiderScheduler
|
||||
from scrapyd.config import Config
|
||||
from scrapyd.scheduler import SpiderScheduler
|
||||
from scrapyd.utils import get_spider_queues
|
||||
|
||||
class SpiderSchedulerTest(unittest.TestCase):
|
||||
|
||||
def setUp(self):
|
||||
d = self.mktemp()
|
||||
eggs_dir = self.eggs_dir = os.path.join(d, 'eggs')
|
||||
dbs_dir = os.path.join(d, 'dbs')
|
||||
os.mkdir(d)
|
||||
os.makedirs(eggs_dir)
|
||||
os.makedirs(dbs_dir)
|
||||
os.makedirs(os.path.join(eggs_dir, 'mybot1'))
|
||||
os.makedirs(os.path.join(eggs_dir, 'mybot2'))
|
||||
config = Config(values={'eggs_dir': eggs_dir, 'dbs_dir': dbs_dir})
|
||||
self.queues = get_spider_queues(config)
|
||||
self.sched = SpiderScheduler(config)
|
||||
|
||||
def test_interface(self):
|
||||
verifyObject(ISpiderScheduler, self.sched)
|
||||
|
||||
def test_list_update_projects(self):
|
||||
self.assertEqual(sorted(self.sched.list_projects()), sorted(['mybot1', 'mybot2']))
|
||||
os.makedirs(os.path.join(self.eggs_dir, 'mybot3'))
|
||||
self.sched.update_projects()
|
||||
self.assertEqual(sorted(self.sched.list_projects()), sorted(['mybot1', 'mybot2', 'mybot3']))
|
||||
|
||||
def test_schedule(self):
|
||||
q = self.queues['mybot1']
|
||||
self.failIf(q.count())
|
||||
self.sched.schedule('mybot1', 'myspider1', a='b')
|
||||
self.sched.schedule('mybot2', 'myspider2', c='d')
|
||||
self.assertEqual(q.pop(), {'name': 'myspider1', 'a': 'b'})
|
||||
q = self.queues['mybot2']
|
||||
self.assertEqual(q.pop(), {'name': 'myspider2', 'c': 'd'})
|
||||
|
||||
|
|
@ -1,66 +0,0 @@
|
|||
from twisted.internet.defer import inlineCallbacks, maybeDeferred
|
||||
from twisted.trial import unittest
|
||||
|
||||
from zope.interface.verify import verifyObject
|
||||
|
||||
from scrapyd.interfaces import ISpiderQueue
|
||||
from scrapyd.spiderqueue import SqliteSpiderQueue
|
||||
|
||||
class SpiderQueueTest(unittest.TestCase):
|
||||
"""This test case can be used easily for testing other SpiderQueue's by
|
||||
just changing the _get_queue() method. It also supports queues with
|
||||
deferred methods.
|
||||
"""
|
||||
|
||||
def setUp(self):
|
||||
self.q = self._get_queue()
|
||||
self.name = 'spider1'
|
||||
self.args = {'arg1': 'val1', 'arg2': 2}
|
||||
self.msg = self.args.copy()
|
||||
self.msg['name'] = self.name
|
||||
|
||||
def _get_queue(self):
|
||||
return SqliteSpiderQueue(':memory:')
|
||||
|
||||
def test_interface(self):
|
||||
verifyObject(ISpiderQueue, self.q)
|
||||
|
||||
@inlineCallbacks
|
||||
def test_add_pop_count(self):
|
||||
c = yield maybeDeferred(self.q.count)
|
||||
self.assertEqual(c, 0)
|
||||
|
||||
yield maybeDeferred(self.q.add, self.name, **self.args)
|
||||
|
||||
c = yield maybeDeferred(self.q.count)
|
||||
self.assertEqual(c, 1)
|
||||
|
||||
m = yield maybeDeferred(self.q.pop)
|
||||
self.assertEqual(m, self.msg)
|
||||
|
||||
c = yield maybeDeferred(self.q.count)
|
||||
self.assertEqual(c, 0)
|
||||
|
||||
@inlineCallbacks
|
||||
def test_list(self):
|
||||
l = yield maybeDeferred(self.q.list)
|
||||
self.assertEqual(l, [])
|
||||
|
||||
yield maybeDeferred(self.q.add, self.name, **self.args)
|
||||
yield maybeDeferred(self.q.add, self.name, **self.args)
|
||||
|
||||
l = yield maybeDeferred(self.q.list)
|
||||
self.assertEqual(l, [self.msg, self.msg])
|
||||
|
||||
@inlineCallbacks
|
||||
def test_clear(self):
|
||||
yield maybeDeferred(self.q.add, self.name, **self.args)
|
||||
yield maybeDeferred(self.q.add, self.name, **self.args)
|
||||
|
||||
c = yield maybeDeferred(self.q.count)
|
||||
self.assertEqual(c, 2)
|
||||
|
||||
yield maybeDeferred(self.q.clear)
|
||||
|
||||
c = yield maybeDeferred(self.q.count)
|
||||
self.assertEqual(c, 0)
|
||||
|
|
@ -1,173 +0,0 @@
|
|||
import unittest
|
||||
from datetime import datetime
|
||||
from decimal import Decimal
|
||||
|
||||
from scrapy.http import Request
|
||||
from scrapyd.sqlite import SqlitePriorityQueue, JsonSqlitePriorityQueue, \
|
||||
PickleSqlitePriorityQueue, SqliteDict, JsonSqliteDict, PickleSqliteDict
|
||||
|
||||
|
||||
class SqliteDictTest(unittest.TestCase):
|
||||
|
||||
dict_class = SqliteDict
|
||||
test_dict = {'hello': 'world', 'int': 1, 'float': 1.5}
|
||||
|
||||
def test_basic_types(self):
|
||||
test = self.test_dict
|
||||
d = self.dict_class()
|
||||
d.update(test)
|
||||
self.failUnlessEqual(d.items(), test.items())
|
||||
d.clear()
|
||||
self.failIf(d.items())
|
||||
|
||||
def test_in(self):
|
||||
d = self.dict_class()
|
||||
self.assertFalse('test' in d)
|
||||
d['test'] = 123
|
||||
self.assertTrue('test' in d)
|
||||
|
||||
def test_keyerror(self):
|
||||
d = self.dict_class()
|
||||
self.assertRaises(KeyError, d.__getitem__, 'test')
|
||||
|
||||
def test_replace(self):
|
||||
d = self.dict_class()
|
||||
self.assertEqual(d.get('test'), None)
|
||||
d['test'] = 123
|
||||
self.assertEqual(d.get('test'), 123)
|
||||
d['test'] = 456
|
||||
self.assertEqual(d.get('test'), 456)
|
||||
|
||||
|
||||
class JsonSqliteDictTest(SqliteDictTest):
|
||||
|
||||
dict_class = JsonSqliteDict
|
||||
test_dict = SqliteDictTest.test_dict.copy()
|
||||
test_dict.update({'list': ['a', 'world'], 'dict': {'some': 'dict'}})
|
||||
|
||||
|
||||
class PickleSqliteDictTest(JsonSqliteDictTest):
|
||||
|
||||
dict_class = PickleSqliteDict
|
||||
test_dict = JsonSqliteDictTest.test_dict.copy()
|
||||
test_dict.update({'decimal': Decimal("10"), 'datetime': datetime.now()})
|
||||
|
||||
def test_request_persistance(self):
|
||||
r1 = Request("http://www.example.com", body="some")
|
||||
d = self.dict_class()
|
||||
d['request'] = r1
|
||||
r2 = d['request']
|
||||
self.failUnless(isinstance(r2, Request))
|
||||
self.failUnlessEqual(r1.url, r2.url)
|
||||
self.failUnlessEqual(r1.body, r2.body)
|
||||
|
||||
|
||||
class SqlitePriorityQueueTest(unittest.TestCase):
|
||||
|
||||
queue_class = SqlitePriorityQueue
|
||||
|
||||
supported_values = ["bytes", u"\xa3", 123, 1.2, True]
|
||||
|
||||
def setUp(self):
|
||||
self.q = self.queue_class()
|
||||
|
||||
def test_empty(self):
|
||||
self.failUnless(self.q.pop() is None)
|
||||
|
||||
def test_one(self):
|
||||
msg = "a message"
|
||||
self.q.put(msg)
|
||||
self.failIf("_id" in msg)
|
||||
self.failUnlessEqual(self.q.pop(), msg)
|
||||
self.failUnless(self.q.pop() is None)
|
||||
|
||||
def test_multiple(self):
|
||||
msg1 = "first message"
|
||||
msg2 = "second message"
|
||||
self.q.put(msg1)
|
||||
self.q.put(msg2)
|
||||
out = []
|
||||
out.append(self.q.pop())
|
||||
out.append(self.q.pop())
|
||||
self.failUnless(msg1 in out)
|
||||
self.failUnless(msg2 in out)
|
||||
self.failUnless(self.q.pop() is None)
|
||||
|
||||
def test_priority(self):
|
||||
msg1 = "message 1"
|
||||
msg2 = "message 2"
|
||||
msg3 = "message 3"
|
||||
msg4 = "message 4"
|
||||
self.q.put(msg1, priority=1.0)
|
||||
self.q.put(msg2, priority=5.0)
|
||||
self.q.put(msg3, priority=3.0)
|
||||
self.q.put(msg4, priority=2.0)
|
||||
self.failUnlessEqual(self.q.pop(), msg2)
|
||||
self.failUnlessEqual(self.q.pop(), msg3)
|
||||
self.failUnlessEqual(self.q.pop(), msg4)
|
||||
self.failUnlessEqual(self.q.pop(), msg1)
|
||||
|
||||
def test_iter_len_clear(self):
|
||||
self.failUnlessEqual(len(self.q), 0)
|
||||
self.failUnlessEqual(list(self.q), [])
|
||||
msg1 = "message 1"
|
||||
msg2 = "message 2"
|
||||
msg3 = "message 3"
|
||||
msg4 = "message 4"
|
||||
self.q.put(msg1, priority=1.0)
|
||||
self.q.put(msg2, priority=5.0)
|
||||
self.q.put(msg3, priority=3.0)
|
||||
self.q.put(msg4, priority=2.0)
|
||||
self.failUnlessEqual(len(self.q), 4)
|
||||
self.failUnlessEqual(list(self.q), \
|
||||
[(msg2, 5.0), (msg3, 3.0), (msg4, 2.0), (msg1, 1.0)])
|
||||
self.q.clear()
|
||||
self.failUnlessEqual(len(self.q), 0)
|
||||
self.failUnlessEqual(list(self.q), [])
|
||||
|
||||
def test_remove(self):
|
||||
self.failUnlessEqual(len(self.q), 0)
|
||||
self.failUnlessEqual(list(self.q), [])
|
||||
msg1 = "good message 1"
|
||||
msg2 = "bad message 2"
|
||||
msg3 = "good message 3"
|
||||
msg4 = "bad message 4"
|
||||
self.q.put(msg1)
|
||||
self.q.put(msg2)
|
||||
self.q.put(msg3)
|
||||
self.q.put(msg4)
|
||||
self.q.remove(lambda x: x.startswith("bad"))
|
||||
self.failUnlessEqual(list(self.q), [(msg1, 0.0), (msg3, 0.0)])
|
||||
|
||||
def test_types(self):
|
||||
for x in self.supported_values:
|
||||
self.q.put(x)
|
||||
self.failUnlessEqual(self.q.pop(), x)
|
||||
|
||||
|
||||
class JsonSqlitePriorityQueueTest(SqlitePriorityQueueTest):
|
||||
|
||||
queue_class = JsonSqlitePriorityQueue
|
||||
|
||||
supported_values = SqlitePriorityQueueTest.supported_values + [
|
||||
["a", "list", 1],
|
||||
{"a": "dict"},
|
||||
]
|
||||
|
||||
|
||||
class PickleSqlitePriorityQueueTest(JsonSqlitePriorityQueueTest):
|
||||
|
||||
queue_class = PickleSqlitePriorityQueue
|
||||
|
||||
supported_values = JsonSqlitePriorityQueueTest.supported_values + [
|
||||
Decimal("10"),
|
||||
datetime.now(),
|
||||
]
|
||||
|
||||
def test_request_persistance(self):
|
||||
r1 = Request("http://www.example.com", body="some")
|
||||
self.q.put(r1)
|
||||
r2 = self.q.pop()
|
||||
self.failUnless(isinstance(r2, Request))
|
||||
self.failUnlessEqual(r1.url, r2.url)
|
||||
self.failUnlessEqual(r1.body, r2.body)
|
||||
|
|
@ -1,51 +0,0 @@
|
|||
import os
|
||||
from pkgutil import get_data
|
||||
from cStringIO import StringIO
|
||||
|
||||
from twisted.trial import unittest
|
||||
|
||||
from scrapy.utils.test import get_pythonpath
|
||||
from scrapyd.interfaces import IEggStorage
|
||||
from scrapyd.utils import get_crawl_args, get_spider_list
|
||||
from scrapyd import get_application
|
||||
|
||||
class UtilsTest(unittest.TestCase):
|
||||
|
||||
def test_get_crawl_args(self):
|
||||
msg = {'_project': 'lolo', '_spider': 'lala'}
|
||||
self.assertEqual(get_crawl_args(msg), ['lala'])
|
||||
msg = {'_project': 'lolo', '_spider': 'lala', 'arg1': u'val1'}
|
||||
cargs = get_crawl_args(msg)
|
||||
self.assertEqual(cargs, ['lala', '-a', 'arg1=val1'])
|
||||
assert all(isinstance(x, str) for x in cargs), cargs
|
||||
|
||||
def test_get_crawl_args_with_settings(self):
|
||||
msg = {'_project': 'lolo', '_spider': 'lala', 'arg1': u'val1', 'settings': {'ONE': 'two'}}
|
||||
cargs = get_crawl_args(msg)
|
||||
self.assertEqual(cargs, ['lala', '-a', 'arg1=val1', '-s', 'ONE=two'])
|
||||
assert all(isinstance(x, str) for x in cargs), cargs
|
||||
|
||||
class GetSpiderListTest(unittest.TestCase):
|
||||
|
||||
def test_get_spider_list(self):
|
||||
path = os.path.abspath(self.mktemp())
|
||||
j = os.path.join
|
||||
eggs_dir = j(path, 'eggs')
|
||||
os.makedirs(eggs_dir)
|
||||
dbs_dir = j(path, 'dbs')
|
||||
os.makedirs(dbs_dir)
|
||||
logs_dir = j(path, 'logs')
|
||||
os.makedirs(logs_dir)
|
||||
os.chdir(path)
|
||||
with open('scrapyd.conf', 'w') as f:
|
||||
f.write("[scrapyd]\n")
|
||||
f.write("eggs_dir = %s\n" % eggs_dir)
|
||||
f.write("dbs_dir = %s\n" % dbs_dir)
|
||||
f.write("logs_dir = %s\n" % logs_dir)
|
||||
app = get_application()
|
||||
eggstorage = app.getComponent(IEggStorage)
|
||||
eggfile = StringIO(get_data(__package__, 'mybot.egg'))
|
||||
eggstorage.put(eggfile, 'mybot', 'r1')
|
||||
spiders = get_spider_list('mybot', pythonpath=get_pythonpath())
|
||||
self.assertEqual(sorted(spiders), ['spider1', 'spider2'])
|
||||
|
||||
|
|
@ -1,67 +0,0 @@
|
|||
import sys
|
||||
import os
|
||||
from subprocess import Popen, PIPE
|
||||
from ConfigParser import NoSectionError
|
||||
|
||||
from scrapyd.spiderqueue import SqliteSpiderQueue
|
||||
from scrapy.utils.python import stringify_dict, unicode_to_str
|
||||
from scrapyd.config import Config
|
||||
|
||||
def get_spider_queues(config):
|
||||
"""Return a dict of Spider Quees keyed by project name"""
|
||||
dbsdir = config.get('dbs_dir', 'dbs')
|
||||
if not os.path.exists(dbsdir):
|
||||
os.makedirs(dbsdir)
|
||||
d = {}
|
||||
for project in get_project_list(config):
|
||||
dbpath = os.path.join(dbsdir, '%s.db' % project)
|
||||
d[project] = SqliteSpiderQueue(dbpath)
|
||||
return d
|
||||
|
||||
def get_project_list(config):
|
||||
"""Get list of projects by inspecting the eggs dir and the ones defined in
|
||||
the scrapyd.conf [settings] section
|
||||
"""
|
||||
eggs_dir = config.get('eggs_dir', 'eggs')
|
||||
if os.path.exists(eggs_dir):
|
||||
projects = os.listdir(eggs_dir)
|
||||
else:
|
||||
projects = []
|
||||
try:
|
||||
projects += [x[0] for x in config.cp.items('settings')]
|
||||
except NoSectionError:
|
||||
pass
|
||||
return projects
|
||||
|
||||
def get_crawl_args(message):
|
||||
"""Return the command-line arguments to use for the scrapy crawl process
|
||||
that will be started for this message
|
||||
"""
|
||||
msg = message.copy()
|
||||
args = [unicode_to_str(msg['_spider'])]
|
||||
del msg['_project'], msg['_spider']
|
||||
settings = msg.pop('settings', {})
|
||||
for k, v in stringify_dict(msg, keys_only=False).items():
|
||||
args += ['-a']
|
||||
args += ['%s=%s' % (k, v)]
|
||||
for k, v in stringify_dict(settings, keys_only=False).items():
|
||||
args += ['-s']
|
||||
args += ['%s=%s' % (k, v)]
|
||||
return args
|
||||
|
||||
def get_spider_list(project, runner=None, pythonpath=None):
|
||||
"""Return the spider list from the given project, using the given runner"""
|
||||
if runner is None:
|
||||
runner = Config().get('runner')
|
||||
env = os.environ.copy()
|
||||
env['SCRAPY_PROJECT'] = project
|
||||
if pythonpath:
|
||||
env['PYTHONPATH'] = pythonpath
|
||||
pargs = [sys.executable, '-m', runner, 'list']
|
||||
proc = Popen(pargs, stdout=PIPE, stderr=PIPE, env=env)
|
||||
out, err = proc.communicate()
|
||||
if proc.returncode:
|
||||
msg = err or out or 'unknown error'
|
||||
raise RuntimeError(msg.splitlines()[-1])
|
||||
return out.splitlines()
|
||||
|
||||
|
|
@ -1,122 +0,0 @@
|
|||
import traceback
|
||||
import uuid
|
||||
from cStringIO import StringIO
|
||||
|
||||
from twisted.python import log
|
||||
|
||||
from scrapy.utils.txweb import JsonResource
|
||||
from .utils import get_spider_list
|
||||
|
||||
class WsResource(JsonResource):
|
||||
|
||||
def __init__(self, root):
|
||||
JsonResource.__init__(self)
|
||||
self.root = root
|
||||
|
||||
def render(self, txrequest):
|
||||
try:
|
||||
return JsonResource.render(self, txrequest)
|
||||
except Exception, e:
|
||||
if self.root.debug:
|
||||
return traceback.format_exc()
|
||||
log.err()
|
||||
r = {"status": "error", "message": str(e)}
|
||||
return self.render_object(r, txrequest)
|
||||
|
||||
class Schedule(WsResource):
|
||||
|
||||
def render_POST(self, txrequest):
|
||||
settings = txrequest.args.pop('setting', [])
|
||||
settings = dict(x.split('=', 1) for x in settings)
|
||||
args = dict((k, v[0]) for k, v in txrequest.args.items())
|
||||
project = args.pop('project')
|
||||
spider = args.pop('spider')
|
||||
args['settings'] = settings
|
||||
jobid = uuid.uuid1().hex
|
||||
args['_job'] = jobid
|
||||
self.root.scheduler.schedule(project, spider, **args)
|
||||
return {"status": "ok", "jobid": jobid}
|
||||
|
||||
class Cancel(WsResource):
|
||||
|
||||
def render_POST(self, txrequest):
|
||||
args = dict((k, v[0]) for k, v in txrequest.args.items())
|
||||
project = args['project']
|
||||
jobid = args['job']
|
||||
signal = args.get('signal', 'TERM')
|
||||
prevstate = None
|
||||
queue = self.root.poller.queues[project]
|
||||
c = queue.remove(lambda x: x["_job"] == jobid)
|
||||
if c:
|
||||
prevstate = "pending"
|
||||
spiders = self.root.launcher.processes.values()
|
||||
for s in spiders:
|
||||
if s.job == jobid:
|
||||
s.transport.signalProcess(signal)
|
||||
prevstate = "running"
|
||||
return {"status": "ok", "prevstate": prevstate}
|
||||
|
||||
class AddVersion(WsResource):
|
||||
|
||||
def render_POST(self, txrequest):
|
||||
project = txrequest.args['project'][0]
|
||||
version = txrequest.args['version'][0]
|
||||
eggf = StringIO(txrequest.args['egg'][0])
|
||||
self.root.eggstorage.put(eggf, project, version)
|
||||
spiders = get_spider_list(project)
|
||||
self.root.update_projects()
|
||||
return {"status": "ok", "project": project, "version": version, \
|
||||
"spiders": len(spiders)}
|
||||
|
||||
class ListProjects(WsResource):
|
||||
|
||||
def render_GET(self, txrequest):
|
||||
projects = self.root.scheduler.list_projects()
|
||||
return {"status": "ok", "projects": projects}
|
||||
|
||||
class ListVersions(WsResource):
|
||||
|
||||
def render_GET(self, txrequest):
|
||||
project = txrequest.args['project'][0]
|
||||
versions = self.root.eggstorage.list(project)
|
||||
return {"status": "ok", "versions": versions}
|
||||
|
||||
class ListSpiders(WsResource):
|
||||
|
||||
def render_GET(self, txrequest):
|
||||
project = txrequest.args['project'][0]
|
||||
spiders = get_spider_list(project, runner=self.root.runner)
|
||||
return {"status": "ok", "spiders": spiders}
|
||||
|
||||
class ListJobs(WsResource):
|
||||
|
||||
def render_GET(self, txrequest):
|
||||
project = txrequest.args['project'][0]
|
||||
spiders = self.root.launcher.processes.values()
|
||||
running = [{"id": s.job, "spider": s.spider} for s in spiders if s.project == project]
|
||||
queue = self.root.poller.queues[project]
|
||||
pending = [{"id": x["_job"], "spider": x["name"]} for x in queue.list()]
|
||||
finished = [{"id": s.job, "spider": s.spider,
|
||||
"start_time": s.start_time.isoformat(' '),
|
||||
"end_time": s.end_time.isoformat(' ')} for s in self.root.launcher.finished
|
||||
if s.project == project]
|
||||
return {"status":"ok", "pending": pending, "running": running, "finished": finished}
|
||||
|
||||
class DeleteProject(WsResource):
|
||||
|
||||
def render_POST(self, txrequest):
|
||||
project = txrequest.args['project'][0]
|
||||
self._delete_version(project)
|
||||
return {"status": "ok"}
|
||||
|
||||
def _delete_version(self, project, version=None):
|
||||
self.root.eggstorage.delete(project, version)
|
||||
self.root.update_projects()
|
||||
|
||||
class DeleteVersion(DeleteProject):
|
||||
|
||||
def render_POST(self, txrequest):
|
||||
project = txrequest.args['project'][0]
|
||||
version = txrequest.args['version'][0]
|
||||
self._delete_version(project, version)
|
||||
return {"status": "ok"}
|
||||
|
|
@ -1,131 +0,0 @@
|
|||
from datetime import datetime
|
||||
|
||||
from twisted.web import resource, static
|
||||
from twisted.application.service import IServiceCollection
|
||||
|
||||
from scrapy.utils.misc import load_object
|
||||
|
||||
from .interfaces import IPoller, IEggStorage, ISpiderScheduler
|
||||
|
||||
class Root(resource.Resource):
|
||||
|
||||
def __init__(self, config, app):
|
||||
resource.Resource.__init__(self)
|
||||
self.debug = config.getboolean('debug', False)
|
||||
self.runner = config.get('runner')
|
||||
logsdir = config.get('logs_dir')
|
||||
itemsdir = config.get('items_dir')
|
||||
self.app = app
|
||||
self.putChild('', Home(self))
|
||||
self.putChild('logs', static.File(logsdir, 'text/plain'))
|
||||
self.putChild('items', static.File(itemsdir, 'text/plain'))
|
||||
self.putChild('jobs', Jobs(self))
|
||||
services = config.items('services', ())
|
||||
for servName, servClsName in services:
|
||||
servCls = load_object(servClsName)
|
||||
self.putChild(servName, servCls(self))
|
||||
self.update_projects()
|
||||
|
||||
def update_projects(self):
|
||||
self.poller.update_projects()
|
||||
self.scheduler.update_projects()
|
||||
|
||||
@property
|
||||
def launcher(self):
|
||||
app = IServiceCollection(self.app, self.app)
|
||||
return app.getServiceNamed('launcher')
|
||||
|
||||
@property
|
||||
def scheduler(self):
|
||||
return self.app.getComponent(ISpiderScheduler)
|
||||
|
||||
@property
|
||||
def eggstorage(self):
|
||||
return self.app.getComponent(IEggStorage)
|
||||
|
||||
@property
|
||||
def poller(self):
|
||||
return self.app.getComponent(IPoller)
|
||||
|
||||
|
||||
class Home(resource.Resource):
|
||||
|
||||
def __init__(self, root):
|
||||
resource.Resource.__init__(self)
|
||||
self.root = root
|
||||
|
||||
def render_GET(self, txrequest):
|
||||
vars = {
|
||||
'projects': ', '.join(self.root.scheduler.list_projects()),
|
||||
}
|
||||
return """
|
||||
<html>
|
||||
<head><title>Scrapyd</title></head>
|
||||
<body>
|
||||
<h1>Scrapyd</h1>
|
||||
<p>Available projects: <b>%(projects)s</b></p>
|
||||
<ul>
|
||||
<li><a href="/jobs">Jobs</a></li>
|
||||
<li><a href="/items/">Items</li>
|
||||
<li><a href="/logs/">Logs</li>
|
||||
<li><a href="http://doc.scrapy.org/en/latest/topics/scrapyd.html">Documentation</a></li>
|
||||
</ul>
|
||||
|
||||
<h2>How to schedule a spider?</h2>
|
||||
|
||||
<p>To schedule a spider you need to use the API (this web UI is only for
|
||||
monitoring)</p>
|
||||
|
||||
<p>Example using <a href="http://curl.haxx.se/">curl</a>:</p>
|
||||
<p><code>curl http://localhost:6800/schedule.json -d project=default -d spider=somespider</code></p>
|
||||
|
||||
<p>For more information about the API, see the <a href="http://doc.scrapy.org/en/latest/topics/scrapyd.html">Scrapyd documentation</a></p>
|
||||
</body>
|
||||
</html>
|
||||
""" % vars
|
||||
|
||||
|
||||
class Jobs(resource.Resource):
|
||||
|
||||
def __init__(self, root):
|
||||
resource.Resource.__init__(self)
|
||||
self.root = root
|
||||
|
||||
def render(self, txrequest):
|
||||
s = "<html><head><title>Scrapyd</title></title>"
|
||||
s += "<body>"
|
||||
s += "<h1>Jobs</h1>"
|
||||
s += "<p><a href='..'>Go back</a></p>"
|
||||
s += "<table border='1'>"
|
||||
s += "<th>Project</th><th>Spider</th><th>Job</th><th>PID</th><th>Runtime</th><th>Log</th><th>Items</th>"
|
||||
s += "<tr><th colspan='7' style='background-color: #ddd'>Pending</th></tr>"
|
||||
for project, queue in self.root.poller.queues.items():
|
||||
for m in queue.list():
|
||||
s += "<tr>"
|
||||
s += "<td>%s</td>" % project
|
||||
s += "<td>%s</td>" % str(m['name'])
|
||||
s += "<td>%s</td>" % str(m['_job'])
|
||||
s += "</tr>"
|
||||
s += "<tr><th colspan='7' style='background-color: #ddd'>Running</th></tr>"
|
||||
for p in self.root.launcher.processes.values():
|
||||
s += "<tr>"
|
||||
for a in ['project', 'spider', 'job', 'pid']:
|
||||
s += "<td>%s</td>" % getattr(p, a)
|
||||
s += "<td>%s</td>" % (datetime.now() - p.start_time)
|
||||
s += "<td><a href='/logs/%s/%s/%s.log'>Log</a></td>" % (p.project, p.spider, p.job)
|
||||
s += "<td><a href='/items/%s/%s/%s.jl'>Items</a></td>" % (p.project, p.spider, p.job)
|
||||
s += "</tr>"
|
||||
s += "<tr><th colspan='7' style='background-color: #ddd'>Finished</th></tr>"
|
||||
for p in self.root.launcher.finished:
|
||||
s += "<tr>"
|
||||
for a in ['project', 'spider', 'job']:
|
||||
s += "<td>%s</td>" % getattr(p, a)
|
||||
s += "<td></td>"
|
||||
s += "<td>%s</td>" % (p.end_time - p.start_time)
|
||||
s += "<td><a href='/logs/%s/%s/%s.log'>Log</a></td>" % (p.project, p.spider, p.job)
|
||||
s += "<td><a href='/items/%s/%s/%s.jl'>Items</a></td>" % (p.project, p.spider, p.job)
|
||||
s += "</tr>"
|
||||
s += "</table>"
|
||||
s += "</body>"
|
||||
s += "</html>"
|
||||
return s
|
||||
6
setup.py
6
setup.py
|
|
@ -1,4 +1,4 @@
|
|||
# Scrapy setup.py script
|
||||
# Scrapy setup.py script
|
||||
#
|
||||
# It doesn't depend on setuptools, but if setuptools is available it'll use
|
||||
# some of its features, like package dependencies.
|
||||
|
|
@ -57,7 +57,7 @@ if root_dir != '':
|
|||
def is_not_module(filename):
|
||||
return os.path.splitext(f)[1] not in ['.py', '.pyc', '.pyo']
|
||||
|
||||
for scrapy_dir in ['scrapy', 'scrapyd']:
|
||||
for scrapy_dir in ['scrapy']:
|
||||
for dirpath, dirnames, filenames in os.walk(scrapy_dir):
|
||||
# Ignore dirnames that start with '.'
|
||||
for i, dirname in enumerate(dirnames):
|
||||
|
|
@ -121,5 +121,5 @@ except ImportError:
|
|||
from distutils.core import setup
|
||||
else:
|
||||
setup_args['install_requires'] = ['Twisted>=8.0', 'w3lib>=1.2', 'lxml', 'pyOpenSSL']
|
||||
|
||||
|
||||
setup(**setup_args)
|
||||
|
|
|
|||
Loading…
Reference in New Issue