Merge branch 'master' into spider.parse

This commit is contained in:
Adrián Chaves 2020-02-06 21:39:09 +01:00 committed by GitHub
commit 24bb9fd5f7
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
68 changed files with 1412 additions and 263 deletions

View File

@ -16,6 +16,8 @@ matrix:
python: 3.5
- env: TOXENV=pinned
python: 3.5
- env: TOXENV=py35-asyncio
python: 3.5.2
- env: TOXENV=py36
python: 3.6
- env: TOXENV=py37
@ -24,6 +26,8 @@ matrix:
python: 3.8
- env: TOXENV=extra-deps
python: 3.8
- env: TOXENV=py38-asyncio
python: 3.8
- env: TOXENV=docs
python: 3.7 # Keep in sync with .readthedocs.yml
install:

View File

@ -34,3 +34,18 @@ def pytest_collection_modifyitems(session, config, items):
items[:] = [item for item in items if isinstance(item, Flake8Item)]
except ImportError:
pass
@pytest.fixture(scope='class')
def reactor_pytest(request):
if not request.cls:
# doctests
return
request.cls.reactor_pytest = request.config.getoption("--reactor")
return request.cls.reactor_pytest
@pytest.fixture(autouse=True)
def only_asyncio(request, reactor_pytest):
if request.node.get_closest_marker('only_asyncio') and reactor_pytest != 'asyncio':
pytest.skip('This test is only run with --reactor=asyncio')

View File

@ -140,7 +140,7 @@ setting the following settings::
While pending requests are below the configured values of
:setting:`CONCURRENT_REQUESTS`, :setting:`CONCURRENT_REQUESTS_PER_DOMAIN` or
:setting:`CONCURRENT_REQUESTS_PER_DOMAIN`, those requests are sent
:setting:`CONCURRENT_REQUESTS_PER_IP`, those requests are sent
concurrently. As a result, the first few requests of a crawl rarely follow the
desired order. Lowering those settings to ``1`` enforces the desired order, but
it significantly slows down the crawl as a whole.
@ -353,7 +353,25 @@ method for this purpose. For example::
for _ in range(item['multiply_by']):
yield deepcopy(item)
Does Scrapy support IPv6 addresses?
-----------------------------------
Yes, by setting :setting:`DNS_RESOLVER` to ``scrapy.resolver.CachingHostnameResolver``.
Note that by doing so, you lose the ability to set a specific timeout for DNS requests
(the value of the :setting:`DNS_TIMEOUT` setting is ignored).
.. _faq-specific-reactor:
How to deal with ``<class 'ValueError'>: filedescriptor out of range in select()`` exceptions?
----------------------------------------------------------------------------------------------
This issue `has been reported`_ to appear when running broad crawls in macOS, where the default
Twisted reactor is :class:`twisted.internet.selectreactor.SelectReactor`. Switching to a
different reactor is possible by using the :setting:`TWISTED_REACTOR` setting.
.. _has been reported: https://github.com/scrapy/scrapy/issues/2905
.. _user agents: https://en.wikipedia.org/wiki/User_agent
.. _LIFO: https://en.wikipedia.org/wiki/Stack_(abstract_data_type)
.. _DFO order: https://en.wikipedia.org/wiki/Depth-first_search

View File

@ -272,7 +272,7 @@ For details, see `Issue #2473 <https://github.com/scrapy/scrapy/issues/2473>`_.
.. _Python: https://www.python.org/
.. _pip: https://pip.pypa.io/en/latest/installing/
.. _lxml: http://lxml.de/
.. _lxml: https://lxml.de/index.html
.. _parsel: https://pypi.python.org/pypi/parsel
.. _w3lib: https://pypi.python.org/pypi/w3lib
.. _twisted: https://twistedmatrix.com/

View File

@ -616,21 +616,25 @@ instance; you still have to yield this Request.
You can also pass a selector to ``response.follow`` instead of a string;
this selector should extract necessary attributes::
for href in response.css('li.next a::attr(href)'):
for href in response.css('ul.pager a::attr(href)'):
yield response.follow(href, callback=self.parse)
For ``<a>`` elements there is a shortcut: ``response.follow`` uses their href
attribute automatically. So the code can be shortened further::
for a in response.css('li.next a'):
for a in response.css('ul.pager a'):
yield response.follow(a, callback=self.parse)
.. note::
To create multiple requests from an iterable, you can use
:meth:`response.follow_all <scrapy.http.TextResponse.follow_all>` instead::
anchors = response.css('ul.pager a')
yield from response.follow_all(anchors, callback=self.parse)
or, shortening it further::
yield from response.follow_all(css='ul.pager a', callback=self.parse)
``response.follow(response.css('li.next a'))`` is not valid because
``response.css`` returns a list-like object with selectors for all results,
not a single selector. A ``for`` loop like in the example above, or
``response.follow(response.css('li.next a')[0])`` is fine.
More examples and patterns
--------------------------
@ -647,13 +651,11 @@ this time for scraping author information::
start_urls = ['http://quotes.toscrape.com/']
def parse(self, response):
# follow links to author pages
for href in response.css('.author + a::attr(href)'):
yield response.follow(href, self.parse_author)
author_page_links = response.css('.author + a')
yield from response.follow_all(author_page_links, self.parse_author)
# follow pagination links
for href in response.css('li.next a::attr(href)'):
yield response.follow(href, self.parse)
pagination_links = response.css('li.next a')
yield from response.follow_all(pagination_links, self.parse)
def parse_author(self, response):
def extract_with_css(query):
@ -669,8 +671,10 @@ This spider will start from the main page, it will follow all the links to the
authors pages calling the ``parse_author`` callback for each of them, and also
the pagination links with the ``parse`` callback as we saw before.
Here we're passing callbacks to ``response.follow`` as positional arguments
to make the code shorter; it also works for ``scrapy.Request``.
Here we're passing callbacks to
:meth:`response.follow_all <scrapy.http.TextResponse.follow_all>` as positional
arguments to make the code shorter; it also works for
:class:`~scrapy.http.Request`.
The ``parse_author`` callback defines a helper function to extract and cleanup the
data from a CSS query and yields the Python dict with the author data.

View File

@ -211,3 +211,10 @@ If your broad crawl shows a high memory usage, in addition to :ref:`crawling in
BFO order <broad-crawls-bfo>` and :ref:`lowering concurrency
<broad-crawls-concurrency>` you should :ref:`debug your memory leaks
<topics-leaks>`.
Install a specific Twisted reactor
==================================
If the crawl is exceeding the system's capabilities, you might want to try
installing a specific Twisted reactor, via the :setting:`TWISTED_REACTOR` setting.

View File

@ -142,7 +142,7 @@ a use case:
Say you want to find the ``Next`` button on the page. Type ``Next`` into the
search bar on the top right of the `Inspector`. You should get two results.
The first is a ``li`` tag with the ``class="text"``, the second the text
The first is a ``li`` tag with the ``class="next"``, the second the text
of an ``a`` tag. Right click on the ``a`` tag and select ``Scroll into View``.
If you hover over the tag, you'll see the button highlighted. From here
we could easily create a :ref:`Link Extractor <topics-link-extractors>` to

View File

@ -199,7 +199,7 @@ CookiesMiddleware
This middleware enables working with sites that require cookies, such as
those that use sessions. It keeps track of cookies sent by web servers, and
send them back on subsequent requests (from that spider), just like web
sends them back on subsequent requests (from that spider), just like web
browsers do.
The following settings can be used to configure the cookie middleware:
@ -672,7 +672,7 @@ sometimes a more nuanced policy is desirable.
This setting still respects ``Cache-Control: no-store`` directives in responses.
If you don't want that, filter ``no-store`` out of the Cache-Control headers in
responses you feedto the cache middleware.
responses you feed to the cache middleware.
.. setting:: HTTPCACHE_IGNORE_RESPONSE_CACHE_CONTROLS
@ -686,7 +686,7 @@ Default: ``[]``
List of Cache-Control directives in responses to be ignored.
Sites often set "no-store", "no-cache", "must-revalidate", etc., but get
upset at the traffic a spider can generate if it respects those
upset at the traffic a spider can generate if it actually respects those
directives. This allows to selectively ignore Cache-Control directives
that are known to be unimportant for the sites being crawled.
@ -868,7 +868,7 @@ Whether the Meta Refresh middleware will be enabled.
METAREFRESH_IGNORE_TAGS
^^^^^^^^^^^^^^^^^^^^^^^
Default: ``['script', 'noscript']``
Default: ``[]``
Meta tags within these tags are ignored.

View File

@ -97,7 +97,6 @@ For Files Pipeline, use::
ITEM_PIPELINES = {'scrapy.pipelines.files.FilesPipeline': 1}
.. note::
You can also use both the Files and Images Pipeline at the same time.
@ -148,6 +147,25 @@ Where:
* ``full`` is a sub-directory to separate full images from thumbnails (if
used). For more info see :ref:`topics-images-thumbnails`.
FTP server storage
------------------
:setting:`FILES_STORE` and :setting:`IMAGES_STORE` can point to an FTP server.
Scrapy will automatically upload the files to the server.
:setting:`FILES_STORE` and :setting:`IMAGES_STORE` should be written in one of the
following forms::
ftp://username:password@address:port/path
ftp://address:port/path
If ``username`` and ``password`` are not provided, they are taken from the :setting:`FTP_USER` and
:setting:`FTP_PASSWORD` settings respectively.
FTP supports two different connection modes: active or passive. Scrapy uses
the passive connection mode by default. To use the active connection mode instead,
set the :setting:`FEED_STORAGE_FTP_ACTIVE` setting to ``True``.
Amazon S3 storage
-----------------
@ -578,4 +596,12 @@ above::
item['image_paths'] = image_paths
return item
To enable your custom media pipeline component you must add its class import path to the
:setting:`ITEM_PIPELINES` setting, like in the following example::
ITEM_PIPELINES = {
'myproject.pipelines.MyImagesPipeline': 300
}
.. _MD5 hash: https://en.wikipedia.org/wiki/MD5

View File

@ -701,6 +701,8 @@ Response objects
.. automethod:: Response.follow
.. automethod:: Response.follow_all
.. _urlparse.urljoin: https://docs.python.org/2/library/urlparse.html#urlparse.urljoin
@ -790,6 +792,8 @@ TextResponse objects
.. automethod:: TextResponse.follow
.. automethod:: TextResponse.follow_all
.. method:: TextResponse.body_as_unicode()
The same as :attr:`text`, but available as a method. This method is

View File

@ -376,6 +376,19 @@ Default: ``10000``
DNS in-memory cache size.
.. setting:: DNS_RESOLVER
DNS_RESOLVER
------------
Default: ``'scrapy.resolver.CachingThreadedResolver'``
The class to be used to resolve DNS names. The default ``scrapy.resolver.CachingThreadedResolver``
supports specifying a timeout for DNS requests via the :setting:`DNS_TIMEOUT` setting,
but works only with IPv4 addresses. Scrapy provides an alternative resolver,
``scrapy.resolver.CachingHostnameResolver``, which supports IPv4/IPv6 addresses but does not
take the :setting:`DNS_TIMEOUT` setting into account.
.. setting:: DNS_TIMEOUT
DNS_TIMEOUT
@ -1241,6 +1254,17 @@ Type of priority queue used by the scheduler. Another available type is
domains in parallel. But currently ``scrapy.pqueues.DownloaderAwarePriorityQueue``
does not work together with :setting:`CONCURRENT_REQUESTS_PER_IP`.
.. setting:: SCRAPER_SLOT_MAX_ACTIVE_SIZE
SCRAPER_SLOT_MAX_ACTIVE_SIZE
----------------------------
Default: ``5_000_000``
Soft limit (in bytes) for response data being processed.
While the sum of the sizes of all responses being processed is above this value,
Scrapy does not process new requests.
.. setting:: SPIDER_CONTRACTS
SPIDER_CONTRACTS
@ -1418,6 +1442,30 @@ command.
The project name must not conflict with the name of custom files or directories
in the ``project`` subdirectory.
.. setting:: TWISTED_REACTOR
TWISTED_REACTOR
---------------
Default: ``None``
Import path of a given Twisted reactor, for instance:
:class:`twisted.internet.asyncioreactor.AsyncioSelectorReactor`.
Scrapy will install this reactor if no other is installed yet, such as when
the ``scrapy`` CLI program is invoked or when using the
:class:`~scrapy.crawler.CrawlerProcess` class. If you are using the
:class:`~scrapy.crawler.CrawlerRunner` class, you need to install the correct
reactor manually. An exception will be raised if the installation fails.
The default value for this option is currently ``None``, which means that Scrapy
will not attempt to install any specific reactor, and the default one defined by
Twisted for the current platform will be used. This is to maintain backward
compatibility and avoid possible problems caused by using a non-default reactor.
For additional information, please see
:doc:`core/howto/choosing-reactor`.
.. setting:: URLLENGTH_LIMIT

View File

@ -18,7 +18,10 @@ addopts =
--ignore=docs/topics/telnetconsole.rst
--ignore=docs/utils
twisted = 1
markers =
only_asyncio: marks tests as only enabled when --reactor=asyncio is passed
flake8-ignore =
W503
# Files that are only meant to provide top-level imports are expected not
# to use any of their imports:
scrapy/core/downloader/handlers/http.py F401
@ -93,6 +96,7 @@ flake8-ignore =
scrapy/loader/__init__.py E501 E128
scrapy/loader/processors.py E501
# scrapy/pipelines
scrapy/pipelines/__init__.py E501
scrapy/pipelines/files.py E116 E501 E266
scrapy/pipelines/images.py E265 E501
scrapy/pipelines/media.py E125 E501 E266
@ -106,7 +110,7 @@ flake8-ignore =
# scrapy/spidermiddlewares
scrapy/spidermiddlewares/httperror.py E501
scrapy/spidermiddlewares/offsite.py E501
scrapy/spidermiddlewares/referer.py E501 E129 W503 W504
scrapy/spidermiddlewares/referer.py E501 E129 W504
scrapy/spidermiddlewares/urllength.py E501
# scrapy/spiders
scrapy/spiders/__init__.py E501 E402
@ -114,6 +118,7 @@ flake8-ignore =
scrapy/spiders/feed.py E501
scrapy/spiders/sitemap.py E501
# scrapy/utils
scrapy/utils/asyncio.py E501
scrapy/utils/benchserver.py E501
scrapy/utils/conf.py E402 E501
scrapy/utils/console.py E306 E305
@ -125,13 +130,13 @@ flake8-ignore =
scrapy/utils/http.py F403 E226
scrapy/utils/httpobj.py E501
scrapy/utils/iterators.py E501 E701
scrapy/utils/log.py E128 W503
scrapy/utils/log.py E128 E501
scrapy/utils/markup.py F403
scrapy/utils/misc.py E501 E226
scrapy/utils/multipart.py F403
scrapy/utils/project.py E501
scrapy/utils/python.py E501
scrapy/utils/reactor.py E226
scrapy/utils/reactor.py E226 E501
scrapy/utils/reqser.py E501
scrapy/utils/request.py E127 E501
scrapy/utils/response.py E501 E128
@ -155,6 +160,7 @@ flake8-ignore =
scrapy/mail.py E402 E128 E501 E502
scrapy/middleware.py E128 E501
scrapy/pqueues.py E501
scrapy/resolver.py E501
scrapy/responsetypes.py E128 E501 E305
scrapy/robotstxt.py E501
scrapy/shell.py E501
@ -223,6 +229,7 @@ flake8-ignore =
tests/test_spidermiddleware_output_chain.py E501 E226
tests/test_spidermiddleware_referer.py E501 F841 E125 E201 E124 E501 E241 E121
tests/test_squeues.py E501 E701 E741
tests/test_utils_asyncio.py E501
tests/test_utils_conf.py E501 E303 E128
tests/test_utils_curl.py E501
tests/test_utils_datatypes.py E402 E501 E305

View File

@ -6,8 +6,8 @@ See documentation in docs/topics/shell.rst
from threading import Thread
from scrapy.commands import ScrapyCommand
from scrapy.shell import Shell
from scrapy.http import Request
from scrapy.shell import Shell
from scrapy.utils.spider import spidercls_for_request, DefaultSpider
from scrapy.utils.url import guess_scheme

View File

@ -4,17 +4,17 @@ import logging
from twisted.internet import defer
from scrapy.exceptions import NotSupported, NotConfigured
from scrapy.utils.httpobj import urlparse_cached
from scrapy.utils.misc import load_object
from scrapy.utils.python import without_none_values
from scrapy import signals
from scrapy.exceptions import NotConfigured, NotSupported
from scrapy.utils.httpobj import urlparse_cached
from scrapy.utils.misc import create_instance, load_object
from scrapy.utils.python import without_none_values
logger = logging.getLogger(__name__)
class DownloadHandlers(object):
class DownloadHandlers:
def __init__(self, crawler):
self._crawler = crawler
@ -49,7 +49,11 @@ class DownloadHandlers(object):
dhcls = load_object(path)
if skip_lazy and getattr(dhcls, 'lazy', True):
return None
dh = dhcls(self._crawler.settings)
dh = create_instance(
objcls=dhcls,
settings=self._crawler.settings,
crawler=self._crawler,
)
except NotConfigured as ex:
self._notconfigured[scheme] = str(ex)
return None

View File

@ -5,12 +5,9 @@ from scrapy.responsetypes import responsetypes
from scrapy.utils.decorators import defers
class DataURIDownloadHandler(object):
class DataURIDownloadHandler:
lazy = False
def __init__(self, settings):
super(DataURIDownloadHandler, self).__init__()
@defers
def download_request(self, request, spider):
uri = parse_data_uri(request.url)

View File

@ -1,14 +1,12 @@
from w3lib.url import file_uri_to_path
from scrapy.responsetypes import responsetypes
from scrapy.utils.decorators import defers
class FileDownloadHandler(object):
class FileDownloadHandler:
lazy = False
def __init__(self, settings):
pass
@defers
def download_request(self, request, spider):
filepath = file_uri_to_path(request.url)

View File

@ -33,8 +33,8 @@ from io import BytesIO
from urllib.parse import unquote
from twisted.internet import reactor
from twisted.protocols.ftp import FTPClient, CommandFailed
from twisted.internet.protocol import Protocol, ClientCreator
from twisted.internet.protocol import ClientCreator, Protocol
from twisted.protocols.ftp import CommandFailed, FTPClient
from scrapy.http import Response
from scrapy.responsetypes import responsetypes
@ -59,10 +59,11 @@ class ReceivedDataProtocol(Protocol):
def close(self):
self.body.close() if self.filename else self.body.seek(0)
_CODE_RE = re.compile(r"\d+")
class FTPDownloadHandler(object):
class FTPDownloadHandler:
lazy = False
CODE_MAPPING = {
@ -75,6 +76,10 @@ class FTPDownloadHandler(object):
self.default_password = settings['FTP_PASSWORD']
self.passive_mode = settings['FTP_PASSIVE_MODE']
@classmethod
def from_crawler(cls, crawler):
return cls(crawler.settings)
def download_request(self, request, spider):
parsed_url = urlparse_cached(request)
user = request.meta.get("ftp_user", self.default_user)

View File

@ -1,17 +1,23 @@
"""Download handlers for http and https schemes
"""
from twisted.internet import reactor
from scrapy.utils.misc import load_object, create_instance
from scrapy.utils.misc import create_instance, load_object
from scrapy.utils.python import to_unicode
class HTTP10DownloadHandler(object):
class HTTP10DownloadHandler:
lazy = False
def __init__(self, settings):
def __init__(self, settings, crawler=None):
self.HTTPClientFactory = load_object(settings['DOWNLOADER_HTTPCLIENTFACTORY'])
self.ClientContextFactory = load_object(settings['DOWNLOADER_CLIENTCONTEXTFACTORY'])
self._settings = settings
self._crawler = crawler
@classmethod
def from_crawler(cls, crawler):
return cls(crawler.settings, crawler)
def download_request(self, request, spider):
"""Return a deferred for the HTTP download"""
@ -22,7 +28,11 @@ class HTTP10DownloadHandler(object):
def _connect(self, factory):
host, port = to_unicode(factory.host), factory.port
if factory.scheme == b'https':
client_context_factory = create_instance(self.ClientContextFactory, settings=self._settings, crawler=None)
client_context_factory = create_instance(
objcls=self.ClientContextFactory,
settings=self._settings,
crawler=self._crawler,
)
return reactor.connectSSL(host, port, factory, client_context_factory)
else:
return reactor.connectTCP(host, port, factory)

View File

@ -1,37 +1,37 @@
"""Download handlers for http and https schemes"""
import re
import logging
import re
import warnings
from io import BytesIO
from time import time
from urllib.parse import urldefrag
from zope.interface import implementer
from twisted.internet import defer, reactor, protocol
from twisted.internet import defer, protocol, reactor
from twisted.internet.endpoints import TCP4ClientEndpoint
from twisted.internet.error import TimeoutError
from twisted.web.client import Agent, HTTPConnectionPool, ResponseDone, ResponseFailed, URI
from twisted.web.http import _DataLoss, PotentialDataLoss
from twisted.web.http_headers import Headers as TxHeaders
from twisted.web.iweb import IBodyProducer, UNKNOWN_LENGTH
from twisted.internet.error import TimeoutError
from twisted.web.http import _DataLoss, PotentialDataLoss
from twisted.web.client import Agent, ResponseDone, HTTPConnectionPool, ResponseFailed, URI
from twisted.internet.endpoints import TCP4ClientEndpoint
from zope.interface import implementer
from scrapy.core.downloader.tls import openssl_methods
from scrapy.core.downloader.webclient import _parse
from scrapy.exceptions import ScrapyDeprecationWarning
from scrapy.http import Headers
from scrapy.responsetypes import responsetypes
from scrapy.core.downloader.webclient import _parse
from scrapy.core.downloader.tls import openssl_methods
from scrapy.utils.misc import load_object, create_instance
from scrapy.utils.misc import create_instance, load_object
from scrapy.utils.python import to_bytes, to_unicode
logger = logging.getLogger(__name__)
class HTTP11DownloadHandler(object):
class HTTP11DownloadHandler:
lazy = False
def __init__(self, settings):
def __init__(self, settings, crawler=None):
self._pool = HTTPConnectionPool(reactor, persistent=True)
self._pool.maxPersistentPerHost = settings.getint('CONCURRENT_REQUESTS_PER_DOMAIN')
self._pool._factory.noisy = False
@ -41,17 +41,17 @@ class HTTP11DownloadHandler(object):
# try method-aware context factory
try:
self._contextFactory = create_instance(
self._contextFactoryClass,
objcls=self._contextFactoryClass,
settings=settings,
crawler=None,
crawler=crawler,
method=self._sslMethod,
)
except TypeError:
# use context factory defaults
self._contextFactory = create_instance(
self._contextFactoryClass,
objcls=self._contextFactoryClass,
settings=settings,
crawler=None,
crawler=crawler,
)
msg = """
'%s' does not accept `method` argument (type OpenSSL.SSL method,\
@ -64,6 +64,10 @@ class HTTP11DownloadHandler(object):
self._fail_on_dataloss = settings.getbool('DOWNLOAD_FAIL_ON_DATALOSS')
self._disconnect_timeout = 1
@classmethod
def from_crawler(cls, crawler):
return cls(crawler.settings, crawler)
def download_request(self, request, spider):
"""Return a deferred for the HTTP download"""
agent = ScrapyAgent(

View File

@ -1,9 +1,10 @@
from urllib.parse import unquote
from scrapy.exceptions import NotConfigured
from scrapy.utils.httpobj import urlparse_cached
from scrapy.utils.boto import is_botocore
from scrapy.core.downloader.handlers.http import HTTPDownloadHandler
from scrapy.exceptions import NotConfigured
from scrapy.utils.boto import is_botocore
from scrapy.utils.httpobj import urlparse_cached
from scrapy.utils.misc import create_instance
def _get_boto_connection():
@ -30,11 +31,12 @@ def _get_boto_connection():
return _S3Connection
class S3DownloadHandler(object):
class S3DownloadHandler:
def __init__(self, settings, aws_access_key_id=None, aws_secret_access_key=None,
def __init__(self, settings, *,
crawler=None,
aws_access_key_id=None, aws_secret_access_key=None,
httpdownloadhandler=HTTPDownloadHandler, **kw):
if not aws_access_key_id:
aws_access_key_id = settings['AWS_ACCESS_KEY_ID']
if not aws_secret_access_key:
@ -67,7 +69,16 @@ class S3DownloadHandler(object):
except Exception as ex:
raise NotConfigured(str(ex))
self._download_http = httpdownloadhandler(settings).download_request
_http_handler = create_instance(
objcls=httpdownloadhandler,
settings=settings,
crawler=crawler,
)
self._download_http = _http_handler.download_request
@classmethod
def from_crawler(cls, crawler, **kwargs):
return cls(crawler.settings, crawler=crawler, **kwargs)
def download_request(self, request, spider):
p = urlparse_cached(request)

View File

@ -8,7 +8,7 @@ from twisted.internet import defer
from scrapy.exceptions import _InvalidOutput
from scrapy.http import Request, Response
from scrapy.middleware import MiddlewareManager
from scrapy.utils.defer import mustbe_deferred
from scrapy.utils.defer import mustbe_deferred, deferred_from_coro
from scrapy.utils.conf import build_component_list
@ -33,7 +33,7 @@ class DownloaderMiddlewareManager(MiddlewareManager):
@defer.inlineCallbacks
def process_request(request):
for method in self.methods['process_request']:
response = yield method(request=request, spider=spider)
response = yield deferred_from_coro(method(request=request, spider=spider))
if response is not None and not isinstance(response, (Response, Request)):
raise _InvalidOutput('Middleware %s.process_request must return None, Response or Request, got %s' % \
(method.__self__.__class__.__name__, response.__class__.__name__))
@ -48,7 +48,7 @@ class DownloaderMiddlewareManager(MiddlewareManager):
defer.returnValue(response)
for method in self.methods['process_response']:
response = yield method(request=request, response=response, spider=spider)
response = yield deferred_from_coro(method(request=request, response=response, spider=spider))
if not isinstance(response, (Response, Request)):
raise _InvalidOutput('Middleware %s.process_response must return Response or Request, got %s' % \
(method.__self__.__class__.__name__, type(response)))
@ -60,7 +60,7 @@ class DownloaderMiddlewareManager(MiddlewareManager):
def process_exception(_failure):
exception = _failure.value
for method in self.methods['process_exception']:
response = yield method(request=request, exception=exception, spider=spider)
response = yield deferred_from_coro(method(request=request, exception=exception, spider=spider))
if response is not None and not isinstance(response, (Response, Request)):
raise _InvalidOutput('Middleware %s.process_exception must return None, Response or Request, got %s' % \
(method.__self__.__class__.__name__, type(response)))

View File

@ -78,7 +78,7 @@ class Scraper(object):
@defer.inlineCallbacks
def open_spider(self, spider):
"""Open the given spider for scraping and allocate resources for it"""
self.slot = Slot()
self.slot = Slot(self.crawler.settings.getint('SCRAPER_SLOT_MAX_ACTIVE_SIZE'))
yield self.itemproc.open_spider(spider)
def close_spider(self, spider):

View File

@ -3,34 +3,36 @@ import pprint
import signal
import warnings
from twisted.internet import reactor, defer
from zope.interface.verify import verifyClass, DoesNotImplement
from twisted.internet import defer
from zope.interface.verify import DoesNotImplement, verifyClass
from scrapy import Spider
from scrapy import signals, Spider
from scrapy.core.engine import ExecutionEngine
from scrapy.resolver import CachingThreadedResolver
from scrapy.interfaces import ISpiderLoader
from scrapy.exceptions import ScrapyDeprecationWarning
from scrapy.extension import ExtensionManager
from scrapy.interfaces import ISpiderLoader
from scrapy.settings import overridden_settings, Settings
from scrapy.signalmanager import SignalManager
from scrapy.exceptions import ScrapyDeprecationWarning
from scrapy.utils.ossignal import install_shutdown_handlers, signal_names
from scrapy.utils.misc import load_object
from scrapy.utils.log import (
LogCounterHandler, configure_logging, log_scrapy_info,
get_scrapy_root_handler, install_scrapy_root_handler)
from scrapy import signals
configure_logging,
get_scrapy_root_handler,
install_scrapy_root_handler,
log_scrapy_info,
LogCounterHandler,
)
from scrapy.utils.misc import create_instance, load_object
from scrapy.utils.ossignal import install_shutdown_handlers, signal_names
from scrapy.utils.reactor import install_reactor, verify_installed_reactor
logger = logging.getLogger(__name__)
class Crawler(object):
class Crawler:
def __init__(self, spidercls, settings=None):
if isinstance(spidercls, Spider):
raise ValueError(
'The spidercls argument must be a class, not an object')
raise ValueError('The spidercls argument must be a class, not an object')
if isinstance(settings, dict) or settings is None:
settings = Settings(settings)
@ -109,7 +111,7 @@ class Crawler(object):
yield defer.maybeDeferred(self.engine.stop)
class CrawlerRunner(object):
class CrawlerRunner:
"""
This is a convenient helper class that keeps track of, manages and runs
crawlers inside an already setup :mod:`~twisted.internet.reactor`.
@ -136,6 +138,7 @@ class CrawlerRunner(object):
self._crawlers = set()
self._active = set()
self.bootstrap_failed = False
self._handle_twisted_reactor()
@property
def spiders(self):
@ -229,6 +232,10 @@ class CrawlerRunner(object):
while self._active:
yield defer.DeferredList(self._active)
def _handle_twisted_reactor(self):
if self.settings.get("TWISTED_REACTOR"):
verify_installed_reactor(self.settings["TWISTED_REACTOR"])
class CrawlerProcess(CrawlerRunner):
"""
@ -261,6 +268,7 @@ class CrawlerProcess(CrawlerRunner):
log_scrapy_info(self.settings)
def _signal_shutdown(self, signum, _):
from twisted.internet import reactor
install_shutdown_handlers(self._signal_kill)
signame = signal_names[signum]
logger.info("Received %(signame)s, shutting down gracefully. Send again to force ",
@ -268,6 +276,7 @@ class CrawlerProcess(CrawlerRunner):
reactor.callFromThread(self._graceful_stop_reactor)
def _signal_kill(self, signum, _):
from twisted.internet import reactor
install_shutdown_handlers(signal.SIG_IGN)
signame = signal_names[signum]
logger.info('Received %(signame)s twice, forcing unclean shutdown',
@ -286,6 +295,7 @@ class CrawlerProcess(CrawlerRunner):
:param boolean stop_after_crawl: stop or not the reactor when all
crawlers have finished
"""
from twisted.internet import reactor
if stop_after_crawl:
d = self.join()
# Don't start the reactor if the deferreds are already fired
@ -293,34 +303,31 @@ class CrawlerProcess(CrawlerRunner):
return
d.addBoth(self._stop_reactor)
reactor.installResolver(self._get_dns_resolver())
resolver_class = load_object(self.settings["DNS_RESOLVER"])
resolver = create_instance(resolver_class, self.settings, self, reactor=reactor)
resolver.install_on_reactor()
tp = reactor.getThreadPool()
tp.adjustPoolsize(maxthreads=self.settings.getint('REACTOR_THREADPOOL_MAXSIZE'))
reactor.addSystemEventTrigger('before', 'shutdown', self.stop)
reactor.run(installSignalHandlers=False) # blocking call
def _get_dns_resolver(self):
if self.settings.getbool('DNSCACHE_ENABLED'):
cache_size = self.settings.getint('DNSCACHE_SIZE')
else:
cache_size = 0
return CachingThreadedResolver(
reactor=reactor,
cache_size=cache_size,
timeout=self.settings.getfloat('DNS_TIMEOUT')
)
def _graceful_stop_reactor(self):
d = self.stop()
d.addBoth(self._stop_reactor)
return d
def _stop_reactor(self, _=None):
from twisted.internet import reactor
try:
reactor.stop()
except RuntimeError: # raised if already stopped or in shutdown stage
pass
def _handle_twisted_reactor(self):
if self.settings.get("TWISTED_REACTOR"):
install_reactor(self.settings["TWISTED_REACTOR"])
super()._handle_twisted_reactor()
def _get_spider_loader(settings):
""" Get SpiderLoader instance from settings """

View File

@ -26,7 +26,7 @@ class HttpCompressionMiddleware(object):
def process_request(self, request, spider):
request.headers.setdefault('Accept-Encoding',
b",".join(ACCEPTED_ENCODINGS))
b", ".join(ACCEPTED_ENCODINGS))
def process_response(self, request, response, spider):

View File

@ -84,6 +84,6 @@ class RetryMiddleware(object):
return retryreq
else:
stats.inc_value('retry/max_reached')
logger.debug("Gave up retrying %(request)s (failed %(retries)d times): %(reason)s",
logger.error("Gave up retrying %(request)s (failed %(retries)d times): %(reason)s",
{'request': request, 'retries': retries, 'reason': reason},
extra={'spider': spider})

View File

@ -23,8 +23,9 @@ __all__ = ['BaseItemExporter', 'PprintItemExporter', 'PickleItemExporter',
class BaseItemExporter(object):
def __init__(self, **kwargs):
self._configure(kwargs)
def __init__(self, dont_fail=False, **kwargs):
self._kwargs = kwargs
self._configure(kwargs, dont_fail=dont_fail)
def _configure(self, options, dont_fail=False):
"""Configure the exporter by poping options from the ``options`` dict.
@ -81,10 +82,10 @@ class BaseItemExporter(object):
class JsonLinesItemExporter(BaseItemExporter):
def __init__(self, file, **kwargs):
self._configure(kwargs, dont_fail=True)
super().__init__(dont_fail=True, **kwargs)
self.file = file
kwargs.setdefault('ensure_ascii', not self.encoding)
self.encoder = ScrapyJSONEncoder(**kwargs)
self._kwargs.setdefault('ensure_ascii', not self.encoding)
self.encoder = ScrapyJSONEncoder(**self._kwargs)
def export_item(self, item):
itemdict = dict(self._get_serialized_fields(item))
@ -95,15 +96,15 @@ class JsonLinesItemExporter(BaseItemExporter):
class JsonItemExporter(BaseItemExporter):
def __init__(self, file, **kwargs):
self._configure(kwargs, dont_fail=True)
super().__init__(dont_fail=True, **kwargs)
self.file = file
# there is a small difference between the behaviour or JsonItemExporter.indent
# and ScrapyJSONEncoder.indent. ScrapyJSONEncoder.indent=None is needed to prevent
# the addition of newlines everywhere
json_indent = self.indent if self.indent is not None and self.indent > 0 else None
kwargs.setdefault('indent', json_indent)
kwargs.setdefault('ensure_ascii', not self.encoding)
self.encoder = ScrapyJSONEncoder(**kwargs)
self._kwargs.setdefault('indent', json_indent)
self._kwargs.setdefault('ensure_ascii', not self.encoding)
self.encoder = ScrapyJSONEncoder(**self._kwargs)
self.first_item = True
def _beautify_newline(self):
@ -134,7 +135,7 @@ class XmlItemExporter(BaseItemExporter):
def __init__(self, file, **kwargs):
self.item_element = kwargs.pop('item_element', 'item')
self.root_element = kwargs.pop('root_element', 'items')
self._configure(kwargs)
super().__init__(**kwargs)
if not self.encoding:
self.encoding = 'utf-8'
self.xg = XMLGenerator(file, encoding=self.encoding)
@ -190,7 +191,7 @@ class XmlItemExporter(BaseItemExporter):
class CsvItemExporter(BaseItemExporter):
def __init__(self, file, include_headers_line=True, join_multivalued=',', **kwargs):
self._configure(kwargs, dont_fail=True)
super().__init__(dont_fail=True, **kwargs)
if not self.encoding:
self.encoding = 'utf-8'
self.include_headers_line = include_headers_line
@ -201,7 +202,7 @@ class CsvItemExporter(BaseItemExporter):
encoding=self.encoding,
newline='' # Windows needs this https://github.com/scrapy/scrapy/issues/3034
)
self.csv_writer = csv.writer(self.stream, **kwargs)
self.csv_writer = csv.writer(self.stream, **self._kwargs)
self._headers_not_written = True
self._join_multivalued = join_multivalued
@ -250,7 +251,7 @@ class CsvItemExporter(BaseItemExporter):
class PickleItemExporter(BaseItemExporter):
def __init__(self, file, protocol=2, **kwargs):
self._configure(kwargs)
super().__init__(**kwargs)
self.file = file
self.protocol = protocol
@ -269,7 +270,7 @@ class MarshalItemExporter(BaseItemExporter):
"""
def __init__(self, file, **kwargs):
self._configure(kwargs)
super().__init__(**kwargs)
self.file = file
def export_item(self, item):
@ -279,7 +280,7 @@ class MarshalItemExporter(BaseItemExporter):
class PprintItemExporter(BaseItemExporter):
def __init__(self, file, **kwargs):
self._configure(kwargs)
super().__init__(**kwargs)
self.file = file
def export_item(self, item):

View File

@ -7,18 +7,16 @@ See documentation in docs/topics/feed-exports.rst
import os
import sys
import logging
import posixpath
from tempfile import NamedTemporaryFile
from datetime import datetime
from urllib.parse import urlparse, unquote
from ftplib import FTP
from zope.interface import Interface, implementer
from twisted.internet import defer, threads
from w3lib.url import file_uri_to_path
from scrapy import signals
from scrapy.utils.ftp import ftp_makedirs_cwd
from scrapy.utils.ftp import ftp_store_file
from scrapy.exceptions import NotConfigured
from scrapy.utils.misc import create_instance, load_object
from scrapy.utils.log import failure_to_exc_info
@ -173,16 +171,11 @@ class FTPFeedStorage(BlockingFeedStorage):
)
def _store_in_thread(self, file):
file.seek(0)
ftp = FTP()
ftp.connect(self.host, self.port)
ftp.login(self.username, self.password)
if self.use_active_mode:
ftp.set_pasv(False)
dirname, filename = posixpath.split(self.path)
ftp_makedirs_cwd(ftp, dirname)
ftp.storbinary('STOR %s' % filename, file)
ftp.quit()
ftp_store_file(
path=self.path, file=file, host=self.host,
port=self.port, username=self.username,
password=self.password, use_active_mode=self.use_active_mode
)
class SpiderSlot(object):

View File

@ -31,7 +31,6 @@ class Request(object_ref):
raise TypeError('callback must be a callable, got %s' % type(callback).__name__)
if errback is not None and not callable(errback):
raise TypeError('errback must be a callable, got %s' % type(errback).__name__)
assert callback or not errback, "Cannot use errback without a callback"
self.callback = callback
self.errback = errback

View File

@ -4,14 +4,15 @@ responses in Scrapy.
See documentation in docs/topics/request-response.rst
"""
from typing import Generator
from urllib.parse import urljoin
from scrapy.http.request import Request
from scrapy.exceptions import NotSupported
from scrapy.http.common import obsolete_setter
from scrapy.http.headers import Headers
from scrapy.http.request import Request
from scrapy.link import Link
from scrapy.utils.trackref import object_ref
from scrapy.http.common import obsolete_setter
from scrapy.exceptions import NotSupported
class Response(object_ref):
@ -41,8 +42,8 @@ class Response(object_ref):
if isinstance(url, str):
self._url = url
else:
raise TypeError('%s url must be str, got %s:' % (type(self).__name__,
type(url).__name__))
raise TypeError('%s url must be str, got %s:' %
(type(self).__name__, type(url).__name__))
url = property(_get_url, obsolete_setter(_set_url, 'url'))
@ -123,14 +124,51 @@ class Response(object_ref):
elif url is None:
raise ValueError("url can't be None")
url = self.urljoin(url)
return Request(url, callback,
method=method,
headers=headers,
body=body,
cookies=cookies,
meta=meta,
encoding=encoding,
priority=priority,
dont_filter=dont_filter,
errback=errback,
cb_kwargs=cb_kwargs)
return Request(
url=url,
callback=callback,
method=method,
headers=headers,
body=body,
cookies=cookies,
meta=meta,
encoding=encoding,
priority=priority,
dont_filter=dont_filter,
errback=errback,
cb_kwargs=cb_kwargs,
)
def follow_all(self, urls, callback=None, method='GET', headers=None, body=None,
cookies=None, meta=None, encoding='utf-8', priority=0,
dont_filter=False, errback=None, cb_kwargs=None):
# type: (...) -> Generator[Request, None, None]
"""
Return an iterable of :class:`~.Request` instances to follow all links
in ``urls``. It accepts the same arguments as ``Request.__init__`` method,
but elements of ``urls`` can be relative URLs or :class:`~scrapy.link.Link` objects,
not only absolute URLs.
:class:`~.TextResponse` provides a :meth:`~.TextResponse.follow_all`
method which supports selectors in addition to absolute/relative URLs
and Link objects.
"""
if not hasattr(urls, '__iter__'):
raise TypeError("'urls' argument must be an iterable")
return (
self.follow(
url=url,
callback=callback,
method=method,
headers=headers,
body=body,
cookies=cookies,
meta=meta,
encoding=encoding,
priority=priority,
dont_filter=dont_filter,
errback=errback,
cb_kwargs=cb_kwargs,
)
for url in urls
)

View File

@ -5,17 +5,19 @@ discovering (through HTTP headers) to base Response class.
See documentation in docs/topics/request-response.rst
"""
from contextlib import suppress
from typing import Generator
from urllib.parse import urljoin
import parsel
from w3lib.encoding import html_to_unicode, resolve_encoding, \
html_body_declared_encoding, http_content_type_encoding
from w3lib.encoding import (html_body_declared_encoding, html_to_unicode,
http_content_type_encoding, resolve_encoding)
from w3lib.html import strip_html5_whitespace
from scrapy.http.request import Request
from scrapy.http import Request
from scrapy.http.response import Response
from scrapy.utils.response import get_base_url
from scrapy.utils.python import memoizemethod_noargs, to_unicode
from scrapy.utils.response import get_base_url
class TextResponse(Response):
@ -40,7 +42,7 @@ class TextResponse(Response):
if isinstance(body, str):
if self._encoding is None:
raise TypeError('Cannot convert unicode body - %s has no encoding' %
type(self).__name__)
type(self).__name__)
self._body = body.encode(self._encoding)
else:
super(TextResponse, self)._set_body(body)
@ -86,8 +88,8 @@ class TextResponse(Response):
if self._cached_benc is None:
content_type = to_unicode(self.headers.get(b'Content-Type', b''))
benc, ubody = html_to_unicode(content_type, self.body,
auto_detect_fun=self._auto_detect_fun,
default_encoding=self._DEFAULT_ENCODING)
auto_detect_fun=self._auto_detect_fun,
default_encoding=self._DEFAULT_ENCODING)
self._cached_benc = benc
self._cached_ubody = ubody
return self._cached_benc
@ -126,13 +128,14 @@ class TextResponse(Response):
It accepts the same arguments as ``Request.__init__`` method,
but ``url`` can be not only an absolute URL, but also
* a relative URL;
* a scrapy.link.Link object (e.g. a link extractor result);
* an attribute Selector (not SelectorList) - e.g.
* a relative URL
* a :class:`~scrapy.link.Link` object, e.g. the result of
:ref:`topics-link-extractors`
* a :class:`~scrapy.selector.Selector` object for a ``<link>`` or ``<a>`` element, e.g.
``response.css('a.my_link')[0]``
* an attribute :class:`~scrapy.selector.Selector` (not SelectorList), e.g.
``response.css('a::attr(href)')[0]`` or
``response.xpath('//img/@src')[0]``.
* a Selector for ``<a>`` or ``<link>`` element, e.g.
``response.css('a.my_link')[0]``.
``response.xpath('//img/@src')[0]``
See :ref:`response-follow-example` for usage examples.
"""
@ -141,7 +144,66 @@ class TextResponse(Response):
elif isinstance(url, parsel.SelectorList):
raise ValueError("SelectorList is not supported")
encoding = self.encoding if encoding is None else encoding
return super(TextResponse, self).follow(url, callback,
return super(TextResponse, self).follow(
url=url,
callback=callback,
method=method,
headers=headers,
body=body,
cookies=cookies,
meta=meta,
encoding=encoding,
priority=priority,
dont_filter=dont_filter,
errback=errback,
cb_kwargs=cb_kwargs,
)
def follow_all(self, urls=None, callback=None, method='GET', headers=None, body=None,
cookies=None, meta=None, encoding=None, priority=0,
dont_filter=False, errback=None, cb_kwargs=None,
css=None, xpath=None):
# type: (...) -> Generator[Request, None, None]
"""
A generator that produces :class:`~.Request` instances to follow all
links in ``urls``. It accepts the same arguments as the :class:`~.Request`'s
``__init__`` method, except that each ``urls`` element does not need to be
an absolute URL, it can be any of the following:
* a relative URL
* a :class:`~scrapy.link.Link` object, e.g. the result of
:ref:`topics-link-extractors`
* a :class:`~scrapy.selector.Selector` object for a ``<link>`` or ``<a>`` element, e.g.
``response.css('a.my_link')[0]``
* an attribute :class:`~scrapy.selector.Selector` (not SelectorList), e.g.
``response.css('a::attr(href)')[0]`` or
``response.xpath('//img/@src')[0]``
In addition, ``css`` and ``xpath`` arguments are accepted to perform the link extraction
within the ``follow_all`` method (only one of ``urls``, ``css`` and ``xpath`` is accepted).
Note that when passing a ``SelectorList`` as argument for the ``urls`` parameter or
using the ``css`` or ``xpath`` parameters, this method will not produce requests for
selectors from which links cannot be obtained (for instance, anchor tags without an
``href`` attribute)
"""
arg_count = len(list(filter(None, (urls, css, xpath))))
if arg_count != 1:
raise ValueError('Please supply exactly one of the following arguments: urls, css, xpath')
if not urls:
if css:
urls = self.css(css)
if xpath:
urls = self.xpath(xpath)
if isinstance(urls, parsel.SelectorList):
selectors = urls
urls = []
for sel in selectors:
with suppress(_InvalidSelector):
urls.append(_url_from_selector(sel))
return super(TextResponse, self).follow_all(
urls=urls,
callback=callback,
method=method,
headers=headers,
body=body,
@ -155,18 +217,24 @@ class TextResponse(Response):
)
class _InvalidSelector(ValueError):
"""
Raised when a URL cannot be obtained from a Selector
"""
def _url_from_selector(sel):
# type: (parsel.Selector) -> str
if isinstance(sel.root, str):
# e.g. ::attr(href) result
return strip_html5_whitespace(sel.root)
if not hasattr(sel.root, 'tag'):
raise ValueError("Unsupported selector: %s" % sel)
raise _InvalidSelector("Unsupported selector: %s" % sel)
if sel.root.tag not in ('a', 'link'):
raise ValueError("Only <a> and <link> elements are supported; got <%s>" %
sel.root.tag)
raise _InvalidSelector("Only <a> and <link> elements are supported; got <%s>" %
sel.root.tag)
href = sel.root.get('href')
if href is None:
raise ValueError("<%s> element has no href attribute: %s" %
(sel.root.tag, sel))
raise _InvalidSelector("<%s> element has no href attribute: %s" %
(sel.root.tag, sel))
return strip_html5_whitespace(href)

View File

@ -6,6 +6,7 @@ See documentation in docs/item-pipeline.rst
from scrapy.middleware import MiddlewareManager
from scrapy.utils.conf import build_component_list
from scrapy.utils.defer import deferred_f_from_coro_f
class ItemPipelineManager(MiddlewareManager):
@ -19,7 +20,7 @@ class ItemPipelineManager(MiddlewareManager):
def _add_middleware(self, pipe):
super(ItemPipelineManager, self)._add_middleware(pipe)
if hasattr(pipe, 'process_item'):
self.methods['process_item'].append(pipe.process_item)
self.methods['process_item'].append(deferred_f_from_coro_f(pipe.process_item))
def process_item(self, item, spider):
return self._process_chain('process_item', item, spider)

View File

@ -11,6 +11,7 @@ import os
import time
from collections import defaultdict
from email.utils import parsedate_tz, mktime_tz
from ftplib import FTP
from io import BytesIO
from urllib.parse import urlparse
@ -26,6 +27,7 @@ from scrapy.utils.python import to_bytes
from scrapy.utils.request import referer_str
from scrapy.utils.boto import is_botocore
from scrapy.utils.datatypes import CaselessDict
from scrapy.utils.ftp import ftp_store_file
logger = logging.getLogger(__name__)
@ -257,6 +259,49 @@ class GCSFilesStore(object):
)
class FTPFilesStore(object):
FTP_USERNAME = None
FTP_PASSWORD = None
USE_ACTIVE_MODE = None
def __init__(self, uri):
assert uri.startswith('ftp://')
u = urlparse(uri)
self.port = u.port
self.host = u.hostname
self.port = int(u.port or 21)
self.username = u.username or self.FTP_USERNAME
self.password = u.password or self.FTP_PASSWORD
self.basedir = u.path.rstrip('/')
def persist_file(self, path, buf, info, meta=None, headers=None):
path = '%s/%s' % (self.basedir, path)
return threads.deferToThread(
ftp_store_file, path=path, file=buf,
host=self.host, port=self.port, username=self.username,
password=self.password, use_active_mode=self.USE_ACTIVE_MODE
)
def stat_file(self, path, info):
def _stat_file(path):
try:
ftp = FTP()
ftp.connect(self.host, self.port)
ftp.login(self.username, self.password)
if self.USE_ACTIVE_MODE:
ftp.set_pasv(False)
file_path = "%s/%s" % (self.basedir, path)
last_modified = float(ftp.voidcmd("MDTM %s" % file_path)[4:].strip())
m = hashlib.md5()
ftp.retrbinary('RETR %s' % file_path, m.update)
return {'last_modified': last_modified, 'checksum': m.hexdigest()}
# The file doesn't exist
except Exception:
return {}
return threads.deferToThread(_stat_file, path)
class FilesPipeline(MediaPipeline):
"""Abstract pipeline that implement the file downloading
@ -283,6 +328,7 @@ class FilesPipeline(MediaPipeline):
'file': FSFilesStore,
's3': S3FilesStore,
'gs': GCSFilesStore,
'ftp': FTPFilesStore
}
DEFAULT_FILES_URLS_FIELD = 'file_urls'
DEFAULT_FILES_RESULT_FIELD = 'files'
@ -330,6 +376,11 @@ class FilesPipeline(MediaPipeline):
gcs_store.GCS_PROJECT_ID = settings['GCS_PROJECT_ID']
gcs_store.POLICY = settings['FILES_STORE_GCS_ACL'] or None
ftp_store = cls.STORE_SCHEMES['ftp']
ftp_store.FTP_USERNAME = settings['FTP_USER']
ftp_store.FTP_PASSWORD = settings['FTP_PASSWORD']
ftp_store.USE_ACTIVE_MODE = settings.getbool('FEED_STORAGE_FTP_ACTIVE')
store_uri = settings['FILES_STORE']
return cls(store_uri, settings=settings)

View File

@ -94,6 +94,11 @@ class ImagesPipeline(FilesPipeline):
gcs_store.GCS_PROJECT_ID = settings['GCS_PROJECT_ID']
gcs_store.POLICY = settings['IMAGES_STORE_GCS_ACL'] or None
ftp_store = cls.STORE_SCHEMES['ftp']
ftp_store.FTP_USERNAME = settings['FTP_USER']
ftp_store.FTP_PASSWORD = settings['FTP_PASSWORD']
ftp_store.USE_ACTIVE_MODE = settings.getbool('FEED_STORAGE_FTP_ACTIVE')
store_uri = settings['IMAGES_STORE']
return cls(store_uri, settings=settings)

View File

@ -1,19 +1,37 @@
from twisted.internet import defer
from twisted.internet.base import ThreadedResolver
from twisted.internet.interfaces import IHostnameResolver, IResolutionReceiver, IResolverSimple
from zope.interface.declarations import implementer, provider
from scrapy.utils.datatypes import LocalCache
# TODO: cache misses
# TODO: cache misses
dnscache = LocalCache(10000)
@implementer(IResolverSimple)
class CachingThreadedResolver(ThreadedResolver):
"""
Default caching resolver. IPv4 only, supports setting a timeout value for DNS requests.
"""
def __init__(self, reactor, cache_size, timeout):
super(CachingThreadedResolver, self).__init__(reactor)
dnscache.limit = cache_size
self.timeout = timeout
@classmethod
def from_crawler(cls, crawler, reactor):
if crawler.settings.getbool('DNSCACHE_ENABLED'):
cache_size = crawler.settings.getint('DNSCACHE_SIZE')
else:
cache_size = 0
return cls(reactor, cache_size, crawler.settings.getfloat('DNS_TIMEOUT'))
def install_on_reactor(self,):
self.reactor.installResolver(self)
def getHostByName(self, name, timeout=None):
if name in dnscache:
return defer.succeed(dnscache[name])
@ -30,3 +48,58 @@ class CachingThreadedResolver(ThreadedResolver):
def _cache_result(self, result, name):
dnscache[name] = result
return result
@implementer(IHostnameResolver)
class CachingHostnameResolver:
"""
Experimental caching resolver. Resolves IPv4 and IPv6 addresses,
does not support setting a timeout value for DNS requests.
"""
def __init__(self, reactor, cache_size):
self.reactor = reactor
self.original_resolver = reactor.nameResolver
dnscache.limit = cache_size
@classmethod
def from_crawler(cls, crawler, reactor):
if crawler.settings.getbool('DNSCACHE_ENABLED'):
cache_size = crawler.settings.getint('DNSCACHE_SIZE')
else:
cache_size = 0
return cls(reactor, cache_size)
def install_on_reactor(self):
self.reactor.installNameResolver(self)
def resolveHostName(self, resolutionReceiver, hostName, portNumber=0,
addressTypes=None, transportSemantics='TCP'):
@provider(IResolutionReceiver)
class CachingResolutionReceiver(resolutionReceiver):
def resolutionBegan(self, resolution):
super(CachingResolutionReceiver, self).resolutionBegan(resolution)
self.resolution = resolution
self.resolved = False
def addressResolved(self, address):
super(CachingResolutionReceiver, self).addressResolved(address)
self.resolved = True
def resolutionComplete(self):
super(CachingResolutionReceiver, self).resolutionComplete()
if self.resolved:
dnscache[hostName] = self.resolution
try:
return dnscache[hostName]
except KeyError:
return self.original_resolver.resolveHostName(
CachingResolutionReceiver(),
hostName,
portNumber,
addressTypes,
transportSemantics
)

View File

@ -58,6 +58,7 @@ DEPTH_PRIORITY = 0
DNSCACHE_ENABLED = True
DNSCACHE_SIZE = 10000
DNS_RESOLVER = 'scrapy.resolver.CachingThreadedResolver'
DNS_TIMEOUT = 60
DOWNLOAD_DELAY = 0
@ -222,7 +223,7 @@ MEMUSAGE_NOTIFY_MAIL = []
MEMUSAGE_WARNING_MB = 0
METAREFRESH_ENABLED = True
METAREFRESH_IGNORE_TAGS = ['script', 'noscript']
METAREFRESH_IGNORE_TAGS = []
METAREFRESH_MAXDELAY = 100
NEWSPIDER_MODULE = ''
@ -252,6 +253,8 @@ SCHEDULER_DISK_QUEUE = 'scrapy.squeues.PickleLifoDiskQueue'
SCHEDULER_MEMORY_QUEUE = 'scrapy.squeues.LifoMemoryQueue'
SCHEDULER_PRIORITY_QUEUE = 'scrapy.pqueues.ScrapyPriorityQueue'
SCRAPER_SLOT_MAX_ACTIVE_SIZE = 5000000
SPIDER_LOADER_CLASS = 'scrapy.spiderloader.SpiderLoader'
SPIDER_LOADER_WARN_ONLY = False
@ -286,6 +289,8 @@ TELNETCONSOLE_HOST = '127.0.0.1'
TELNETCONSOLE_USERNAME = 'scrapy'
TELNETCONSOLE_PASSWORD = None
TWISTED_REACTOR = None
SPIDER_CONTRACTS = {}
SPIDER_CONTRACTS_BASE = {
'scrapy.contracts.default.UrlContract': 1,

View File

@ -7,7 +7,7 @@ import os
import signal
import warnings
from twisted.internet import reactor, threads, defer
from twisted.internet import threads, defer
from twisted.python import threadable
from w3lib.url import any_to_uri
@ -98,6 +98,7 @@ class Shell(object):
return spider
def fetch(self, request_or_url, spider=None, redirect=True, **kwargs):
from twisted.internet import reactor
if isinstance(request_or_url, Request):
request = request_or_url
else:

View File

@ -1,11 +1,15 @@
"""
Helper functions for dealing with Twisted deferreds
"""
import asyncio
import inspect
from functools import wraps
from twisted.internet import defer, reactor, task
from twisted.internet import defer, task
from twisted.python import failure
from scrapy.exceptions import IgnoreRequest
from scrapy.utils.reactor import is_asyncio_reactor_installed
def defer_fail(_failure):
@ -15,6 +19,7 @@ def defer_fail(_failure):
It delays by 100ms so reactor has a chance to go through readers and writers
before attending pending delayed calls, so do not set delay to zero.
"""
from twisted.internet import reactor
d = defer.Deferred()
reactor.callLater(0.1, d.errback, _failure)
return d
@ -27,6 +32,7 @@ def defer_succeed(result):
It delays by 100ms so reactor has a chance to go trough readers and writers
before attending pending delayed calls, so do not set delay to zero.
"""
from twisted.internet import reactor
d = defer.Deferred()
reactor.callLater(0.1, d.callback, result)
return d
@ -113,3 +119,37 @@ def iter_errback(iterable, errback, *a, **kw):
break
except Exception:
errback(failure.Failure(), *a, **kw)
def _isfuture(o):
# workaround for Python before 3.5.3 not having asyncio.isfuture
if hasattr(asyncio, 'isfuture'):
return asyncio.isfuture(o)
return isinstance(o, asyncio.Future)
def deferred_from_coro(o):
"""Converts a coroutine into a Deferred, or returns the object as is if it isn't a coroutine"""
if isinstance(o, defer.Deferred):
return o
if _isfuture(o) or inspect.isawaitable(o):
if not is_asyncio_reactor_installed():
# wrapping the coroutine directly into a Deferred, this doesn't work correctly with coroutines
# that use asyncio, e.g. "await asyncio.sleep(1)"
return defer.ensureDeferred(o)
else:
# wrapping the coroutine into a Future and then into a Deferred, this requires AsyncioSelectorReactor
return defer.Deferred.fromFuture(asyncio.ensure_future(o))
return o
def deferred_f_from_coro_f(coro_f):
""" Converts a coroutine function into a function that returns a Deferred.
The coroutine function will be called at the time when the wrapper is called. Wrapper args will be passed to it.
This is useful for callback chains, as callback functions are called with the previous callback result.
"""
@wraps(coro_f)
def f(*coro_args, **coro_kwargs):
return deferred_from_coro(coro_f(*coro_args, **coro_kwargs))
return f

View File

@ -1,4 +1,6 @@
from ftplib import error_perm
import posixpath
from ftplib import error_perm, FTP
from posixpath import dirname
@ -14,3 +16,20 @@ def ftp_makedirs_cwd(ftp, path, first_call=True):
ftp.mkd(path)
if first_call:
ftp.cwd(path)
def ftp_store_file(
*, path, file, host, port,
username, password, use_active_mode=False):
"""Opens a FTP connection with passed credentials,sets current directory
to the directory extracted from given path, then uploads the file to server
"""
with FTP() as ftp:
ftp.connect(host, port)
ftp.login(username, password)
if use_active_mode:
ftp.set_pasv(False)
file.seek(0)
dirname, filename = posixpath.split(path)
ftp_makedirs_cwd(ftp, dirname)
ftp.storbinary('STOR %s' % filename, file)

View File

@ -1,16 +1,16 @@
# -*- coding: utf-8 -*-
import sys
import logging
import sys
import warnings
from logging.config import dictConfig
from twisted.python.failure import Failure
from twisted.python import log as twisted_log
from twisted.python.failure import Failure
import scrapy
from scrapy.settings import Settings
from scrapy.exceptions import ScrapyDeprecationWarning
from scrapy.settings import Settings
from scrapy.utils.versions import scrapy_components_versions
@ -148,6 +148,8 @@ def log_scrapy_info(settings):
{'versions': ", ".join("%s %s" % (name, version)
for name, version in scrapy_components_versions()
if name != "Scrapy")})
from twisted.internet import reactor
logger.debug("Using reactor: %s.%s", reactor.__module__, reactor.__class__.__name__)
class StreamLogger(object):

View File

@ -1,7 +1,5 @@
import signal
from twisted.internet import reactor
signal_names = {}
for signame in dir(signal):
@ -17,6 +15,7 @@ def install_shutdown_handlers(function, override_sigint=True):
SIGINT handler won't be install if there is already a handler in place
(e.g. Pdb)
"""
from twisted.internet import reactor
reactor._handleSignals()
signal.signal(signal.SIGTERM, function)
if signal.getsignal(signal.SIGINT) == signal.default_int_handler or \

View File

@ -1,8 +1,14 @@
from twisted.internet import reactor, error
import asyncio
from contextlib import suppress
from twisted.internet import asyncioreactor, error
from scrapy.utils.misc import load_object
def listen_tcp(portrange, host, factory):
"""Like reactor.listenTCP but tries different ports in a range."""
from twisted.internet import reactor
assert len(portrange) <= 2, "invalid portrange: %s" % portrange
if not portrange:
return reactor.listenTCP(0, factory, interface=host)
@ -30,6 +36,7 @@ class CallLaterOnce(object):
self._call = None
def schedule(self, delay=0):
from twisted.internet import reactor
if self._call is None:
self._call = reactor.callLater(delay, self)
@ -40,3 +47,31 @@ class CallLaterOnce(object):
def __call__(self):
self._call = None
return self._func(*self._a, **self._kw)
def install_reactor(reactor_path):
reactor_class = load_object(reactor_path)
if reactor_class is asyncioreactor.AsyncioSelectorReactor:
with suppress(error.ReactorAlreadyInstalledError):
asyncioreactor.install(asyncio.get_event_loop())
else:
*module, _ = reactor_path.split(".")
installer_path = module + ["install"]
installer = load_object(".".join(installer_path))
with suppress(error.ReactorAlreadyInstalledError):
installer()
def verify_installed_reactor(reactor_path):
from twisted.internet import reactor
reactor_class = load_object(reactor_path)
if not isinstance(reactor, reactor_class):
msg = "The installed reactor ({}.{}) does not match the requested one ({})".format(
reactor.__module__, reactor.__class__.__name__, reactor_path
)
raise Exception(msg)
def is_asyncio_reactor_installed():
from twisted.internet import reactor
return isinstance(reactor, asyncioreactor.AsyncioSelectorReactor)

View File

@ -2,14 +2,15 @@ import logging
import inspect
from scrapy.spiders import Spider
from scrapy.utils.misc import arg_to_iter
from scrapy.utils.defer import deferred_from_coro
from scrapy.utils.misc import arg_to_iter
logger = logging.getLogger(__name__)
def iterate_spider_output(result):
return arg_to_iter(result)
return arg_to_iter(deferred_from_coro(result))
def iter_spider_classes(module):

View File

@ -2,6 +2,9 @@
This module contains some assorted functions used in tests
"""
from __future__ import absolute_import
from posixpath import split
import asyncio
import os
from importlib import import_module
@ -63,6 +66,26 @@ def get_gcs_content_and_delete(bucket, path):
return content, acl, blob
def get_ftp_content_and_delete(
path, host, port, username,
password, use_active_mode=False):
from ftplib import FTP
ftp = FTP()
ftp.connect(host, port)
ftp.login(username, password)
if use_active_mode:
ftp.set_pasv(False)
ftp_data = []
def buffer_data(data):
ftp_data.append(data)
ftp.retrbinary('RETR %s' % path, buffer_data)
dirname, filename = split(path)
ftp.cwd(dirname)
ftp.delete(filename)
return "".join(ftp_data)
def get_crawler(spidercls=None, settings_dict=None):
"""Return an unconfigured Crawler object. If settings_dict is given, it
will be used to populate the crawler settings with a project level
@ -96,3 +119,10 @@ def assert_samelines(testcase, text1, text2, msg=None):
line endings between platforms
"""
testcase.assertEqual(text1.splitlines(), text2.splitlines(), msg)
def get_from_asyncio_queue(value):
q = asyncio.Queue()
getter = q.get()
q.put_nowait(value)
return getter

View File

@ -0,0 +1,15 @@
import scrapy
from scrapy.crawler import CrawlerProcess
class IPv6Spider(scrapy.Spider):
name = "ipv6_spider"
start_urls = ["http://[::1]"]
process = CrawlerProcess(settings={
"RETRY_ENABLED": False,
"DNS_RESOLVER": "scrapy.resolver.CachingHostnameResolver",
})
process.crawl(IPv6Spider)
process.start()

View File

@ -0,0 +1,16 @@
import scrapy
from scrapy.crawler import CrawlerProcess
class NoRequestsSpider(scrapy.Spider):
name = 'no_request'
def start_requests(self):
return []
process = CrawlerProcess(settings={
"TWISTED_REACTOR": "twisted.internet.asyncioreactor.AsyncioSelectorReactor",
})
process.crawl(NoRequestsSpider)
process.start()

View File

@ -0,0 +1,21 @@
import asyncio
from twisted.internet import asyncioreactor
asyncioreactor.install(asyncio.get_event_loop())
import scrapy
from scrapy.crawler import CrawlerProcess
class NoRequestsSpider(scrapy.Spider):
name = 'no_request'
def start_requests(self):
return []
process = CrawlerProcess(settings={
"TWISTED_REACTOR": "twisted.internet.asyncioreactor.AsyncioSelectorReactor",
})
process.crawl(NoRequestsSpider)
process.start()

View File

@ -0,0 +1,12 @@
import scrapy
from scrapy.crawler import CrawlerProcess
class IPv6Spider(scrapy.Spider):
name = "ipv6_spider"
start_urls = ["http://[::1]"]
process = CrawlerProcess(settings={"RETRY_ENABLED": False})
process.crawl(IPv6Spider)
process.start()

View File

@ -0,0 +1,13 @@
import scrapy
from scrapy.crawler import CrawlerProcess
class AsyncioReactorSpider(scrapy.Spider):
name = 'asyncio_reactor'
process = CrawlerProcess(settings={
"TWISTED_REACTOR": "twisted.internet.asyncioreactor.AsyncioSelectorReactor",
})
process.crawl(AsyncioReactorSpider)
process.start()

View File

@ -0,0 +1,13 @@
import scrapy
from scrapy.crawler import CrawlerProcess
class PollReactorSpider(scrapy.Spider):
name = 'poll_reactor'
process = CrawlerProcess(settings={
"TWISTED_REACTOR": "twisted.internet.pollreactor.PollReactor",
})
process.crawl(PollReactorSpider)
process.start()

View File

@ -0,0 +1,13 @@
import scrapy
from scrapy.crawler import CrawlerProcess
class SelectReactorSpider(scrapy.Spider):
name = 'epoll_reactor'
process = CrawlerProcess(settings={
"TWISTED_REACTOR": "twisted.internet.selectreactor.SelectReactor",
})
process.crawl(SelectReactorSpider)
process.start()

View File

@ -4,7 +4,7 @@ mitmproxy; python_version >= '3.6'
mitmproxy<4.0.0; python_version < '3.6'
pytest
pytest-cov
pytest-twisted
pytest-twisted >= 1.11
pytest-xdist
sybil
testfixtures

View File

@ -0,0 +1,25 @@
<html>
<head>
<base href='http://example.com' />
<title>Sample page with anchor tags containing no href attribute, to test the TextResponse.follow_all method</title>
</head>
<body>
<div class="quote">
<span class="text">“The world as we have created it is a process of our
thinking. It cannot be changed without changing our thinking.”</span>
<span>
by <small class="author">Albert Einstein</small>
<a href="/author/Albert-Einstein">(about)</a>
</span>
<div id="pagination" class="pagination">
Tags:
<a href="/page/1/">Page 1</a>
<a>Current</a>
<a href="/page/3/">Page 3</a>
<a href="/page/4/">Page 4</a>
</div>
</div>
</body>
</html>

View File

@ -1,14 +1,18 @@
"""
Some spiders used for testing and benchmarking
"""
import asyncio
import time
from urllib.parse import urlencode
from twisted.internet import defer
from scrapy.http import Request
from scrapy.item import Item
from scrapy.linkextractors import LinkExtractor
from scrapy.spiders import Spider
from scrapy.spiders.crawl import CrawlSpider, Rule
from scrapy.utils.test import get_from_asyncio_queue
class MockServerSpider(Spider):
@ -83,6 +87,36 @@ class SimpleSpider(MetaSpider):
self.logger.info("Got response %d" % response.status)
class AsyncDefSpider(SimpleSpider):
name = 'asyncdef'
async def parse(self, response):
await defer.succeed(42)
self.logger.info("Got response %d" % response.status)
class AsyncDefAsyncioSpider(SimpleSpider):
name = 'asyncdef_asyncio'
async def parse(self, response):
await asyncio.sleep(0.2)
status = await get_from_asyncio_queue(response.status)
self.logger.info("Got response %d" % status)
class AsyncDefAsyncioReturnSpider(SimpleSpider):
name = 'asyncdef_asyncio_return'
async def parse(self, response):
await asyncio.sleep(0.2)
status = await get_from_asyncio_queue(response.status)
self.logger.info("Got response %d" % status)
return [{'id': 1}, {'id': 2}]
class ItemSpider(FollowAllSpider):
name = 'item'

View File

@ -295,6 +295,16 @@ class BadSpider(scrapy.Spider):
self.assertIn("start_requests", log)
self.assertIn("badspider.py", log)
def test_asyncio_enabled_true(self):
log = self.get_log(self.debug_log_spider, args=[
'-s', 'TWISTED_REACTOR=twisted.internet.asyncioreactor.AsyncioSelectorReactor'
])
self.assertIn("Using reactor: twisted.internet.asyncioreactor.AsyncioSelectorReactor", log)
def test_asyncio_enabled_false(self):
log = self.get_log(self.debug_log_spider, args=[])
self.assertNotIn("Using reactor: twisted.internet.asyncioreactor.AsyncioSelectorReactor", log)
class BenchCommandTest(CommandTest):

View File

@ -1,17 +1,29 @@
import json
import logging
from pytest import mark
from testfixtures import LogCapture
from twisted.internet import defer
from twisted.trial.unittest import TestCase
from scrapy import signals
from scrapy.crawler import CrawlerRunner
from scrapy.http import Request
from scrapy.utils.python import to_unicode
from tests.mockserver import MockServer
from tests.spiders import (BrokenStartRequestsSpider, CrawlSpiderWithErrback,
CrawlSpiderWithParseMethod, DelaySpider, SimpleSpider,
DuplicateStartRequestsSpider, FollowAllSpider, SingleRequestSpider)
from tests.spiders import (
AsyncDefAsyncioReturnSpider,
AsyncDefAsyncioSpider,
AsyncDefSpider,
BrokenStartRequestsSpider,
CrawlSpiderWithErrback,
CrawlSpiderWithParseMethod,
DelaySpider,
DuplicateStartRequestsSpider,
FollowAllSpider,
SimpleSpider,
SingleRequestSpider,
)
class CrawlTestCase(TestCase):
@ -333,3 +345,35 @@ class CrawlSpiderTestCase(TestCase):
self.assertIn("[errback] status 404", str(log))
self.assertIn("[errback] status 500", str(log))
self.assertIn("[errback] status 501", str(log))
@defer.inlineCallbacks
def test_async_def_parse(self):
self.runner.crawl(AsyncDefSpider, self.mockserver.url("/status?n=200"), mockserver=self.mockserver)
with LogCapture() as log:
yield self.runner.join()
self.assertIn("Got response 200", str(log))
@mark.only_asyncio()
@defer.inlineCallbacks
def test_async_def_asyncio_parse(self):
runner = CrawlerRunner({"ASYNCIO_REACTOR": True})
runner.crawl(AsyncDefAsyncioSpider, self.mockserver.url("/status?n=200"), mockserver=self.mockserver)
with LogCapture() as log:
yield runner.join()
self.assertIn("Got response 200", str(log))
@mark.only_asyncio()
@defer.inlineCallbacks
def test_async_def_asyncio_parse_list(self):
items = []
def _on_item_scraped(item):
items.append(item)
crawler = self.runner.create_crawler(AsyncDefAsyncioReturnSpider)
crawler.signals.connect(_on_item_scraped, signals.item_scraped)
with LogCapture() as log:
yield crawler.crawl(self.mockserver.url("/status?n=200"), mockserver=self.mockserver)
self.assertIn("Got response 200", str(log))
self.assertIn({'id': 1}, items)
self.assertIn({'id': 2}, items)

View File

@ -4,9 +4,10 @@ import subprocess
import sys
import warnings
from pytest import raises, mark
from testfixtures import LogCapture
from twisted.internet import defer
from twisted.trial import unittest
from pytest import raises
import scrapy
from scrapy.crawler import Crawler, CrawlerRunner, CrawlerProcess
@ -207,6 +208,7 @@ class NoRequestsSpider(scrapy.Spider):
return []
@mark.usefixtures('reactor_pytest')
class CrawlerRunnerHasSpider(unittest.TestCase):
@defer.inlineCallbacks
@ -250,6 +252,41 @@ class CrawlerRunnerHasSpider(unittest.TestCase):
self.assertEqual(runner.bootstrap_failed, True)
def test_crawler_runner_asyncio_enabled_true(self):
if self.reactor_pytest == 'asyncio':
runner = CrawlerRunner(settings={
"TWISTED_REACTOR": "twisted.internet.asyncioreactor.AsyncioSelectorReactor",
})
else:
msg = r"The installed reactor \(.*?\) does not match the requested one \(.*?\)"
with self.assertRaisesRegex(Exception, msg):
runner = CrawlerRunner(settings={
"TWISTED_REACTOR": "twisted.internet.asyncioreactor.AsyncioSelectorReactor",
})
@defer.inlineCallbacks
def test_crawler_process_asyncio_enabled_true(self):
with LogCapture(level=logging.DEBUG) as log:
if self.reactor_pytest == 'asyncio':
runner = CrawlerProcess(settings={
"TWISTED_REACTOR": "twisted.internet.asyncioreactor.AsyncioSelectorReactor",
})
yield runner.crawl(NoRequestsSpider)
self.assertIn("Using reactor: twisted.internet.asyncioreactor.AsyncioSelectorReactor", str(log))
else:
msg = r"The installed reactor \(.*?\) does not match the requested one \(.*?\)"
with self.assertRaisesRegex(Exception, msg):
runner = CrawlerProcess(settings={
"TWISTED_REACTOR": "twisted.internet.asyncioreactor.AsyncioSelectorReactor",
})
@defer.inlineCallbacks
def test_crawler_process_asyncio_enabled_false(self):
runner = CrawlerProcess(settings={"TWISTED_REACTOR": None})
with LogCapture(level=logging.DEBUG) as log:
yield runner.crawl(NoRequestsSpider)
self.assertNotIn("Using reactor: twisted.internet.asyncioreactor.AsyncioSelectorReactor", str(log))
class CrawlerProcessSubprocess(unittest.TestCase):
script_dir = os.path.join(os.path.abspath(os.path.dirname(__file__)), 'CrawlerProcess')
@ -265,3 +302,47 @@ class CrawlerProcessSubprocess(unittest.TestCase):
def test_simple(self):
log = self.run_script('simple.py')
self.assertIn('Spider closed (finished)', log)
self.assertNotIn("Using reactor: twisted.internet.asyncioreactor.AsyncioSelectorReactor", log)
def test_asyncio_enabled_no_reactor(self):
log = self.run_script('asyncio_enabled_no_reactor.py')
self.assertIn('Spider closed (finished)', log)
self.assertIn("Using reactor: twisted.internet.asyncioreactor.AsyncioSelectorReactor", log)
def test_asyncio_enabled_reactor(self):
log = self.run_script('asyncio_enabled_reactor.py')
self.assertIn('Spider closed (finished)', log)
self.assertIn("Using reactor: twisted.internet.asyncioreactor.AsyncioSelectorReactor", log)
def test_ipv6_default_name_resolver(self):
log = self.run_script('default_name_resolver.py')
self.assertIn('Spider closed (finished)', log)
self.assertIn("twisted.internet.error.DNSLookupError: DNS lookup failed: no results for hostname lookup: ::1.", log)
self.assertIn("'downloader/exception_type_count/twisted.internet.error.DNSLookupError': 1,", log)
def test_ipv6_alternative_name_resolver(self):
log = self.run_script('alternative_name_resolver.py')
self.assertIn('Spider closed (finished)', log)
self.assertTrue(any([
"twisted.internet.error.ConnectionRefusedError" in log,
"twisted.internet.error.ConnectError" in log,
]))
self.assertTrue(any([
"'downloader/exception_type_count/twisted.internet.error.ConnectionRefusedError': 1," in log,
"'downloader/exception_type_count/twisted.internet.error.ConnectError': 1," in log,
]))
def test_reactor_select(self):
log = self.run_script("twisted_reactor_select.py")
self.assertIn("Spider closed (finished)", log)
self.assertIn("Using reactor: twisted.internet.selectreactor.SelectReactor", log)
def test_reactor_poll(self):
log = self.run_script("twisted_reactor_poll.py")
self.assertIn("Spider closed (finished)", log)
self.assertIn("Using reactor: twisted.internet.pollreactor.PollReactor", log)
def test_reactor_asyncio(self):
log = self.run_script("twisted_reactor_asyncio.py")
self.assertIn("Spider closed (finished)", log)
self.assertIn("Using reactor: twisted.internet.asyncioreactor.AsyncioSelectorReactor", log)

View File

@ -1,21 +1,20 @@
import contextlib
import os
import shutil
import tempfile
from unittest import mock
import contextlib
from testfixtures import LogCapture
from twisted.trial import unittest
from twisted.cred import checkers, credentials, portal
from twisted.internet import defer, error, reactor
from twisted.protocols.policies import WrappingFactory
from twisted.python.filepath import FilePath
from twisted.internet import reactor, defer, error
from twisted.web import server, static, util, resource
from twisted.trial import unittest
from twisted.web import resource, server, static, util
from twisted.web._newclient import ResponseFailed
from twisted.web.http import _DataLoss
from twisted.web.test.test_webclient import ForeverTakingResource, \
NoLengthResource, HostHeaderResource, \
PayloadResource
from twisted.cred import portal, checkers, credentials
from twisted.web.test.test_webclient import (ForeverTakingResource, HostHeaderResource,
NoLengthResource, PayloadResource)
from w3lib.url import path_to_file_uri
from scrapy.core.downloader.handlers import DownloadHandlers
@ -26,39 +25,38 @@ from scrapy.core.downloader.handlers.http10 import HTTP10DownloadHandler
from scrapy.core.downloader.handlers.http11 import HTTP11DownloadHandler
from scrapy.core.downloader.handlers.s3 import S3DownloadHandler
from scrapy.spiders import Spider
from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning
from scrapy.http import Headers, Request
from scrapy.http.response.text import TextResponse
from scrapy.responsetypes import responsetypes
from scrapy.settings import Settings
from scrapy.utils.test import get_crawler, skip_if_no_boto
from scrapy.spiders import Spider
from scrapy.utils.misc import create_instance
from scrapy.utils.python import to_bytes
from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning
from scrapy.utils.test import get_crawler, skip_if_no_boto
from tests.mockserver import MockServer, ssl_context_factory, Echo
from tests.spiders import SingleRequestSpider
class DummyDH(object):
class DummyDH:
lazy = False
def __init__(self, crawler):
pass
class DummyLazyDH(object):
class DummyLazyDH:
# Default is lazy for backward compatibility
def __init__(self, crawler):
pass
pass
class OffDH(object):
class OffDH:
lazy = False
def __init__(self, crawler):
raise NotConfigured
@classmethod
def from_crawler(cls, crawler):
return cls(crawler)
class LoadTestCase(unittest.TestCase):
@ -106,7 +104,8 @@ class FileTestCase(unittest.TestCase):
self.tmpname = self.mktemp()
with open(self.tmpname + '^', 'w') as f:
f.write('0123456789')
self.download_request = FileDownloadHandler(Settings()).download_request
handler = create_instance(FileDownloadHandler, None, get_crawler())
self.download_request = handler.download_request
def tearDown(self):
os.unlink(self.tmpname + '^')
@ -239,7 +238,7 @@ class HttpTestCase(unittest.TestCase):
else:
self.port = reactor.listenTCP(0, self.wrapper, interface=self.host)
self.portno = self.port.getHost().port
self.download_handler = self.download_handler_cls(Settings())
self.download_handler = create_instance(self.download_handler_cls, None, get_crawler())
self.download_request = self.download_handler.download_request
@defer.inlineCallbacks
@ -479,9 +478,8 @@ class Http11TestCase(HttpTestCase):
return self.test_download_broken_content_allow_data_loss('broken-chunked')
def test_download_broken_content_allow_data_loss_via_setting(self, url='broken'):
download_handler = self.download_handler_cls(Settings({
'DOWNLOAD_FAIL_ON_DATALOSS': False,
}))
crawler = get_crawler(settings_dict={'DOWNLOAD_FAIL_ON_DATALOSS': False})
download_handler = create_instance(self.download_handler_cls, None, crawler)
request = Request(self.getURL(url))
d = download_handler.download_request(request, Spider('foo'))
d.addCallback(lambda r: r.flags)
@ -499,9 +497,8 @@ class Https11TestCase(Http11TestCase):
@defer.inlineCallbacks
def test_tls_logging(self):
download_handler = self.download_handler_cls(Settings({
'DOWNLOADER_CLIENT_TLS_VERBOSE_LOGGING': True,
}))
crawler = get_crawler(settings_dict={'DOWNLOADER_CLIENT_TLS_VERBOSE_LOGGING': True})
download_handler = create_instance(self.download_handler_cls, None, crawler)
try:
with LogCapture() as log_capture:
request = Request(self.getURL('file'))
@ -568,8 +565,8 @@ class Https11CustomCiphers(unittest.TestCase):
0, self.wrapper, ssl_context_factory(self.keyfile, self.certfile, cipher_string='CAMELLIA256-SHA'),
interface=self.host)
self.portno = self.port.getHost().port
self.download_handler = self.download_handler_cls(
Settings({'DOWNLOADER_CLIENT_TLS_CIPHERS': 'CAMELLIA256-SHA'}))
crawler = get_crawler(settings_dict={'DOWNLOADER_CLIENT_TLS_CIPHERS': 'CAMELLIA256-SHA'})
self.download_handler = create_instance(self.download_handler_cls, None, crawler)
self.download_request = self.download_handler.download_request
@defer.inlineCallbacks
@ -665,7 +662,7 @@ class HttpProxyTestCase(unittest.TestCase):
wrapper = WrappingFactory(site)
self.port = reactor.listenTCP(0, wrapper, interface='127.0.0.1')
self.portno = self.port.getHost().port
self.download_handler = self.download_handler_cls(Settings())
self.download_handler = create_instance(self.download_handler_cls, None, get_crawler())
self.download_request = self.download_handler.download_request
@defer.inlineCallbacks
@ -731,9 +728,7 @@ class Http11ProxyTestCase(HttpProxyTestCase):
self.assertIn(domain, timeout.osError)
class HttpDownloadHandlerMock(object):
def __init__(self, settings):
pass
class HttpDownloadHandlerMock:
def download_request(self, request, spider):
return request
@ -743,9 +738,13 @@ class S3AnonTestCase(unittest.TestCase):
def setUp(self):
skip_if_no_boto()
self.s3reqh = S3DownloadHandler(Settings(),
httpdownloadhandler=HttpDownloadHandlerMock,
#anon=True, # is implicit
crawler = get_crawler()
self.s3reqh = create_instance(
objcls=S3DownloadHandler,
settings=None,
crawler=crawler,
httpdownloadhandler=HttpDownloadHandlerMock,
# anon=True, # implicit
)
self.download_request = self.s3reqh.download_request
self.spider = Spider('foo')
@ -771,9 +770,15 @@ class S3TestCase(unittest.TestCase):
def setUp(self):
skip_if_no_boto()
s3reqh = S3DownloadHandler(Settings(), self.AWS_ACCESS_KEY_ID,
self.AWS_SECRET_ACCESS_KEY,
httpdownloadhandler=HttpDownloadHandlerMock)
crawler = get_crawler()
s3reqh = create_instance(
objcls=S3DownloadHandler,
settings=None,
crawler=crawler,
aws_access_key_id=self.AWS_ACCESS_KEY_ID,
aws_secret_access_key=self.AWS_SECRET_ACCESS_KEY,
httpdownloadhandler=HttpDownloadHandlerMock,
)
self.download_request = s3reqh.download_request
self.spider = Spider('foo')
@ -793,7 +798,13 @@ class S3TestCase(unittest.TestCase):
def test_extra_kw(self):
try:
S3DownloadHandler(Settings(), extra_kw=True)
crawler = get_crawler()
create_instance(
objcls=S3DownloadHandler,
settings=None,
crawler=crawler,
extra_kw=True,
)
except Exception as e:
self.assertIsInstance(e, (TypeError, NotConfigured))
else:
@ -935,7 +946,8 @@ class BaseFTPTestCase(unittest.TestCase):
self.factory = FTPFactory(portal=p)
self.port = reactor.listenTCP(0, self.factory, interface="127.0.0.1")
self.portNum = self.port.getHost().port
self.download_handler = FTPDownloadHandler(Settings())
crawler = get_crawler()
self.download_handler = create_instance(FTPDownloadHandler, crawler.settings, crawler)
self.addCleanup(self.port.stopListening)
def tearDown(self):
@ -1049,7 +1061,8 @@ class AnonymousFTPTestCase(BaseFTPTestCase):
userAnonymous=self.username)
self.port = reactor.listenTCP(0, self.factory, interface="127.0.0.1")
self.portNum = self.port.getHost().port
self.download_handler = FTPDownloadHandler(Settings())
crawler = get_crawler()
self.download_handler = create_instance(FTPDownloadHandler, crawler.settings, crawler)
self.addCleanup(self.port.stopListening)
def tearDown(self):
@ -1059,7 +1072,8 @@ class AnonymousFTPTestCase(BaseFTPTestCase):
class DataURITestCase(unittest.TestCase):
def setUp(self):
self.download_handler = DataURIDownloadHandler(Settings())
crawler = get_crawler()
self.download_handler = create_instance(DataURIDownloadHandler, crawler.settings, crawler)
self.download_request = self.download_handler.download_request
self.spider = Spider('foo')

View File

@ -1,5 +1,9 @@
import asyncio
from unittest import mock
from pytest import mark
from twisted.internet import defer
from twisted.internet.defer import Deferred
from twisted.trial.unittest import TestCase
from twisted.python.failure import Failure
@ -7,7 +11,7 @@ from scrapy.http import Request, Response
from scrapy.spiders import Spider
from scrapy.exceptions import _InvalidOutput
from scrapy.core.downloader.middleware import DownloaderMiddlewareManager
from scrapy.utils.test import get_crawler
from scrapy.utils.test import get_crawler, get_from_asyncio_queue
from scrapy.utils.python import to_bytes
@ -177,3 +181,75 @@ class ProcessExceptionInvalidOutput(ManagerTestCase):
dfd.addBoth(results.append)
self.assertIsInstance(results[0], Failure)
self.assertIsInstance(results[0].value, _InvalidOutput)
class MiddlewareUsingDeferreds(ManagerTestCase):
"""Middlewares using Deferreds should work"""
def test_deferred(self):
resp = Response('http://example.com/index.html')
class DeferredMiddleware:
def cb(self, result):
return result
def process_request(self, request, spider):
d = Deferred()
d.addCallback(self.cb)
d.callback(resp)
return d
self.mwman._add_middleware(DeferredMiddleware())
req = Request('http://example.com/index.html')
download_func = mock.MagicMock()
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
self._wait(dfd)
self.assertIs(results[0], resp)
self.assertFalse(download_func.called)
class MiddlewareUsingCoro(ManagerTestCase):
"""Middlewares using asyncio coroutines should work"""
def test_asyncdef(self):
resp = Response('http://example.com/index.html')
class CoroMiddleware:
async def process_request(self, request, spider):
await defer.succeed(42)
return resp
self.mwman._add_middleware(CoroMiddleware())
req = Request('http://example.com/index.html')
download_func = mock.MagicMock()
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
self._wait(dfd)
self.assertIs(results[0], resp)
self.assertFalse(download_func.called)
@mark.only_asyncio()
def test_asyncdef_asyncio(self):
resp = Response('http://example.com/index.html')
class CoroMiddleware:
async def process_request(self, request, spider):
await asyncio.sleep(0.1)
result = await get_from_asyncio_queue(resp)
return result
self.mwman._add_middleware(CoroMiddleware())
req = Request('http://example.com/index.html')
download_func = mock.MagicMock()
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
self._wait(dfd)
self.assertIs(results[0], resp)
self.assertFalse(download_func.called)

View File

@ -48,7 +48,7 @@ class HttpCompressionTest(TestCase):
}
response = Response('http://scrapytest.org/', body=body, headers=headers)
response.request = Request('http://scrapytest.org', headers={'Accept-Encoding': 'gzip,deflate'})
response.request = Request('http://scrapytest.org', headers={'Accept-Encoding': 'gzip, deflate'})
return response
def test_process_request(self):
@ -56,7 +56,7 @@ class HttpCompressionTest(TestCase):
assert 'Accept-Encoding' not in request.headers
self.mw.process_request(request, self.spider)
self.assertEqual(request.headers.get('Accept-Encoding'),
b','.join(ACCEPTED_ENCODINGS))
b', '.join(ACCEPTED_ENCODINGS))
def test_process_response_gzip(self):
response = self._getresponse('gzip')

View File

@ -300,19 +300,21 @@ class MetaRefreshMiddlewareTest(unittest.TestCase):
body = ('''<noscript><meta http-equiv="refresh" '''
'''content="0;URL='http://example.org/newpage'"></noscript>''')
rsp = HtmlResponse(req.url, body=body.encode())
response = self.mw.process_response(req, rsp, self.spider)
assert isinstance(response, Response)
req2 = self.mw.process_response(req, rsp, self.spider)
assert isinstance(req2, Request)
self.assertEqual(req2.url, 'http://example.org/newpage')
def test_ignore_tags_empty_list(self):
crawler = get_crawler(Spider, {'METAREFRESH_IGNORE_TAGS': []})
def test_ignore_tags_1_x_list(self):
"""Test that Scrapy 1.x behavior remains possible"""
settings = {'METAREFRESH_IGNORE_TAGS': ['script', 'noscript']}
crawler = get_crawler(Spider, settings)
mw = MetaRefreshMiddleware.from_crawler(crawler)
req = Request(url='http://example.org')
body = ('''<noscript><meta http-equiv="refresh" '''
'''content="0;URL='http://example.org/newpage'"></noscript>''')
rsp = HtmlResponse(req.url, body=body.encode())
req2 = mw.process_response(req, rsp, self.spider)
assert isinstance(req2, Request)
self.assertEqual(req2.url, 'http://example.org/newpage')
response = mw.process_response(req, rsp, self.spider)
assert isinstance(response, Response)
if __name__ == "__main__":
unittest.main()

View File

@ -244,25 +244,41 @@ class RequestTest(unittest.TestCase):
self.assertRaises(AttributeError, setattr, r, 'url', 'http://example2.com')
self.assertRaises(AttributeError, setattr, r, 'body', 'xxx')
def test_callback_is_callable(self):
def test_callback_and_errback(self):
def a_function():
pass
r = self.request_class('http://example.com')
self.assertIsNone(r.callback)
r = self.request_class('http://example.com', a_function)
self.assertIs(r.callback, a_function)
with self.assertRaises(TypeError):
self.request_class('http://example.com', 'a_function')
def test_errback_is_callable(self):
def a_function():
pass
r = self.request_class('http://example.com')
self.assertIsNone(r.errback)
r = self.request_class('http://example.com', a_function, errback=a_function)
self.assertIs(r.errback, a_function)
r1 = self.request_class('http://example.com')
self.assertIsNone(r1.callback)
self.assertIsNone(r1.errback)
r2 = self.request_class('http://example.com', callback=a_function)
self.assertIs(r2.callback, a_function)
self.assertIsNone(r2.errback)
r3 = self.request_class('http://example.com', errback=a_function)
self.assertIsNone(r3.callback)
self.assertIs(r3.errback, a_function)
r4 = self.request_class(
url='http://example.com',
callback=a_function,
errback=a_function,
)
self.assertIs(r4.callback, a_function)
self.assertIs(r4.errback, a_function)
def test_callback_and_errback_type(self):
with self.assertRaises(TypeError):
self.request_class('http://example.com', a_function, errback='a_function')
self.request_class('http://example.com', callback='a_function')
with self.assertRaises(TypeError):
self.request_class('http://example.com', errback='a_function')
with self.assertRaises(TypeError):
self.request_class(
url='http://example.com',
callback='a_function',
errback='a_function',
)
def test_from_curl(self):
# Note: more curated tests regarding curl conversion are in

View File

@ -141,6 +141,8 @@ class BaseResponseTest(unittest.TestCase):
r.css('body')
r.xpath('//body')
# Response.follow
def test_follow_url_absolute(self):
self._assert_followed_url('http://foo.example.com',
'http://foo.example.com')
@ -164,6 +166,72 @@ class BaseResponseTest(unittest.TestCase):
def test_follow_whitespace_link(self):
self._assert_followed_url(Link('http://example.com/foo '),
'http://example.com/foo%20')
# Response.follow_all
def test_follow_all_absolute(self):
url_list = ['http://example.org', 'http://www.example.org',
'http://example.com', 'http://www.example.com']
self._assert_followed_all_urls(url_list, url_list)
def test_follow_all_relative(self):
relative = ['foo', 'bar', 'foo/bar', 'bar/foo']
absolute = [
'http://example.com/foo',
'http://example.com/bar',
'http://example.com/foo/bar',
'http://example.com/bar/foo',
]
self._assert_followed_all_urls(relative, absolute)
def test_follow_all_links(self):
absolute = [
'http://example.com/foo',
'http://example.com/bar',
'http://example.com/foo/bar',
'http://example.com/bar/foo',
]
links = map(Link, absolute)
self._assert_followed_all_urls(links, absolute)
def test_follow_all_invalid(self):
r = self.response_class("http://example.com")
if self.response_class == Response:
with self.assertRaises(TypeError):
list(r.follow_all(urls=None))
with self.assertRaises(TypeError):
list(r.follow_all(urls=12345))
with self.assertRaises(ValueError):
list(r.follow_all(urls=[None]))
else:
with self.assertRaises(ValueError):
list(r.follow_all(urls=None))
with self.assertRaises(TypeError):
list(r.follow_all(urls=12345))
with self.assertRaises(ValueError):
list(r.follow_all(urls=[None]))
def test_follow_all_whitespace(self):
relative = ['foo ', 'bar ', 'foo/bar ', 'bar/foo ']
absolute = [
'http://example.com/foo%20',
'http://example.com/bar%20',
'http://example.com/foo/bar%20',
'http://example.com/bar/foo%20',
]
self._assert_followed_all_urls(relative, absolute)
def test_follow_all_whitespace_links(self):
absolute = [
'http://example.com/foo ',
'http://example.com/bar ',
'http://example.com/foo/bar ',
'http://example.com/bar/foo ',
]
links = map(Link, absolute)
expected = [u.replace(' ', '%20') for u in absolute]
self._assert_followed_all_urls(links, expected)
def _assert_followed_url(self, follow_obj, target_url, response=None):
if response is None:
response = self._links_response()
@ -171,8 +239,21 @@ class BaseResponseTest(unittest.TestCase):
self.assertEqual(req.url, target_url)
return req
def _assert_followed_all_urls(self, follow_obj, target_urls, response=None):
if response is None:
response = self._links_response()
followed = response.follow_all(follow_obj)
for req, target in zip(followed, target_urls):
self.assertEqual(req.url, target)
yield req
def _links_response(self):
body = get_testdata('link_extractor', 'sgml_linkextractor.html')
body = get_testdata('link_extractor', 'linkextractor.html')
resp = self.response_class('http://example.com/index', body=body)
return resp
def _links_response_no_href(self):
body = get_testdata('link_extractor', 'linkextractor_no_href.html')
resp = self.response_class('http://example.com/index', body=body)
return resp
@ -481,6 +562,53 @@ class TextResponseTest(BaseResponseTest):
)
self.assertEqual(req.encoding, 'cp1251')
def test_follow_all_css(self):
expected = [
'http://example.com/sample3.html',
'http://example.com/innertag.html',
]
response = self._links_response()
extracted = [r.url for r in response.follow_all(css='a[href*="example.com"]')]
self.assertEqual(expected, extracted)
def test_follow_all_css_skip_invalid(self):
expected = [
'http://example.com/page/1/',
'http://example.com/page/3/',
'http://example.com/page/4/',
]
response = self._links_response_no_href()
extracted1 = [r.url for r in response.follow_all(css='.pagination a')]
self.assertEqual(expected, extracted1)
extracted2 = [r.url for r in response.follow_all(response.css('.pagination a'))]
self.assertEqual(expected, extracted2)
def test_follow_all_xpath(self):
expected = [
'http://example.com/sample3.html',
'http://example.com/innertag.html',
]
response = self._links_response()
extracted = response.follow_all(xpath='//a[contains(@href, "example.com")]')
self.assertEqual(expected, [r.url for r in extracted])
def test_follow_all_xpath_skip_invalid(self):
expected = [
'http://example.com/page/1/',
'http://example.com/page/3/',
'http://example.com/page/4/',
]
response = self._links_response_no_href()
extracted1 = [r.url for r in response.follow_all(xpath='//div[@id="pagination"]/a')]
self.assertEqual(expected, extracted1)
extracted2 = [r.url for r in response.follow_all(response.xpath('//div[@id="pagination"]/a'))]
self.assertEqual(expected, extracted2)
def test_follow_all_too_many_arguments(self):
response = self._links_response()
with self.assertRaises(ValueError):
response.follow_all(css='a[href*="example.com"]', xpath='//a[contains(@href, "example.com")]')
class HtmlResponseTest(TextResponseTest):

View File

@ -19,7 +19,7 @@ class Base:
escapes_whitespace = False
def setUp(self):
body = get_testdata('link_extractor', 'sgml_linkextractor.html')
body = get_testdata('link_extractor', 'linkextractor.html')
self.response = HtmlResponse(url='http://example.com/index', body=body)
def test_urls_type(self):

View File

@ -10,12 +10,13 @@ from urllib.parse import urlparse
from twisted.trial import unittest
from twisted.internet import defer
from scrapy.pipelines.files import FilesPipeline, FSFilesStore, S3FilesStore, GCSFilesStore
from scrapy.pipelines.files import FilesPipeline, FSFilesStore, S3FilesStore, GCSFilesStore, FTPFilesStore
from scrapy.item import Item, Field
from scrapy.http import Request, Response
from scrapy.settings import Settings
from scrapy.utils.test import assert_aws_environ, get_s3_content_and_delete
from scrapy.utils.test import assert_gcs_environ, get_gcs_content_and_delete
from scrapy.utils.test import get_ftp_content_and_delete
from scrapy.utils.boto import is_botocore
@ -367,6 +368,31 @@ class TestGCSFilesStore(unittest.TestCase):
self.assertIn(expected_policy, acl)
class TestFTPFileStore(unittest.TestCase):
@defer.inlineCallbacks
def test_persist(self):
uri = os.environ.get('FTP_TEST_FILE_URI')
if not uri:
raise unittest.SkipTest("No FTP URI available for testing")
data = b"TestFTPFilesStore: \xe2\x98\x83"
buf = BytesIO(data)
meta = {'foo': 'bar'}
path = 'full/filename'
store = FTPFilesStore(uri)
empty_dict = yield store.stat_file(path, info=None)
self.assertEqual(empty_dict, {})
yield store.persist_file(path, buf, info=None, meta=meta, headers=None)
stat = yield store.stat_file(path, info=None)
self.assertIn('last_modified', stat)
self.assertIn('checksum', stat)
self.assertEqual(stat['checksum'], 'd113d66b2ec7258724a268bd88eef6b6')
path = '%s/%s' % (store.basedir, path)
content = get_ftp_content_and_delete(
path, store.host, store.port,
store.username, store.password, store.USE_ACTIVE_MODE)
self.assertEqual(data.decode(), content)
class ItemWithFiles(Item):
file_urls = Field()
files = Field()

View File

@ -1,9 +1,12 @@
import asyncio
from pytest import mark
from twisted.internet import defer
from twisted.internet.defer import Deferred
from twisted.trial import unittest
from scrapy import Spider, signals, Request
from scrapy.utils.test import get_crawler
from scrapy.utils.test import get_crawler, get_from_asyncio_queue
from tests.mockserver import MockServer
@ -26,6 +29,20 @@ class DeferredPipeline:
return d
class AsyncDefPipeline:
async def process_item(self, item, spider):
await defer.succeed(42)
item['pipeline_passed'] = True
return item
class AsyncDefAsyncioPipeline:
async def process_item(self, item, spider):
await asyncio.sleep(0.2)
item['pipeline_passed'] = await get_from_asyncio_queue(True)
return item
class ItemSpider(Spider):
name = 'itemspider'
@ -69,3 +86,16 @@ class PipelineTestCase(unittest.TestCase):
crawler = self._create_crawler(DeferredPipeline)
yield crawler.crawl(mockserver=self.mockserver)
self.assertEqual(len(self.items), 1)
@defer.inlineCallbacks
def test_asyncdef_pipeline(self):
crawler = self._create_crawler(AsyncDefPipeline)
yield crawler.crawl(mockserver=self.mockserver)
self.assertEqual(len(self.items), 1)
@mark.only_asyncio()
@defer.inlineCallbacks
def test_asyncdef_asyncio_pipeline(self):
crawler = self._create_crawler(AsyncDefAsyncioPipeline)
yield crawler.crawl(mockserver=self.mockserver)
self.assertEqual(len(self.items), 1)

View File

@ -0,0 +1,17 @@
from unittest import TestCase
from pytest import mark
from scrapy.utils.reactor import is_asyncio_reactor_installed, install_reactor
@mark.usefixtures('reactor_pytest')
class AsyncioTest(TestCase):
def test_is_asyncio_reactor_installed(self):
# the result should depend only on the pytest --reactor argument
self.assertEqual(is_asyncio_reactor_installed(), self.reactor_pytest == 'asyncio')
def test_install_asyncio_reactor(self):
# this should do nothing
install_reactor("twisted.internet.asyncioreactor.AsyncioSelectorReactor")

15
tox.ini
View File

@ -56,7 +56,6 @@ deps =
pyOpenSSL==16.2.0
queuelib==1.4.2
service_identity==16.0.0
six==1.10.0
Twisted==17.9.0
w3lib==1.17.0
zope.interface==4.1.3
@ -96,3 +95,17 @@ changedir = {[docs]changedir}
deps = {[docs]deps}
commands =
sphinx-build -W -b linkcheck . {envtmpdir}/linkcheck
[asyncio]
commands =
{[testenv]commands} --reactor=asyncio
[testenv:py35-asyncio]
basepython = python3.5
deps = {[testenv]deps}
commands = {[asyncio]commands}
[testenv:py38-asyncio]
basepython = python3.8
deps = {[testenv]deps}
commands = {[asyncio]commands}