mirror of https://github.com/scrapy/scrapy.git
Merge remote-tracking branch 'origin/master' into asyncio-parse-asyncgen-proper-rebased
This commit is contained in:
commit
ecfc924ca8
|
|
@ -0,0 +1,19 @@
|
|||
[flake8]
|
||||
|
||||
max-line-length = 119
|
||||
ignore = W503
|
||||
|
||||
exclude =
|
||||
# Exclude files that are meant to provide top-level imports
|
||||
# E402: Module level import not at top of file
|
||||
# F401: Module imported but unused
|
||||
scrapy/__init__.py E402
|
||||
scrapy/core/downloader/handlers/http.py F401
|
||||
scrapy/http/__init__.py F401
|
||||
scrapy/linkextractors/__init__.py E402 F401
|
||||
scrapy/selector/__init__.py F401
|
||||
scrapy/spiders/__init__.py E402 F401
|
||||
|
||||
# Issues pending a review:
|
||||
scrapy/utils/url.py F403 F405
|
||||
tests/test_loader.py E741
|
||||
|
|
@ -14,6 +14,8 @@ htmlcov/
|
|||
.coverage
|
||||
.pytest_cache/
|
||||
.coverage.*
|
||||
coverage.*
|
||||
test-output.*
|
||||
.cache/
|
||||
.mypy_cache/
|
||||
/tests/keys/localhost.crt
|
||||
|
|
|
|||
18
conftest.py
18
conftest.py
|
|
@ -1,6 +1,7 @@
|
|||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
from twisted.web.http import H2_ENABLED
|
||||
|
||||
from scrapy.utils.reactor import install_reactor
|
||||
|
||||
|
|
@ -20,10 +21,19 @@ collect_ignore = [
|
|||
*_py_files("tests/CrawlerRunner"),
|
||||
]
|
||||
|
||||
for line in open('tests/ignores.txt'):
|
||||
file_path = line.strip()
|
||||
if file_path and file_path[0] != '#':
|
||||
collect_ignore.append(file_path)
|
||||
with open('tests/ignores.txt') as reader:
|
||||
for line in reader:
|
||||
file_path = line.strip()
|
||||
if file_path and file_path[0] != '#':
|
||||
collect_ignore.append(file_path)
|
||||
|
||||
if not H2_ENABLED:
|
||||
collect_ignore.extend(
|
||||
(
|
||||
'scrapy/core/downloader/handlers/http2.py',
|
||||
*_py_files("scrapy/core/http2"),
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture()
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ PYTHON = python
|
|||
SPHINXOPTS =
|
||||
PAPER =
|
||||
SOURCES =
|
||||
SHELL = /bin/bash
|
||||
SHELL = /usr/bin/env bash
|
||||
|
||||
ALLSPHINXOPTS = -b $(BUILDER) -d build/doctrees \
|
||||
-D latex_elements.papersize=$(PAPER) \
|
||||
|
|
|
|||
|
|
@ -227,6 +227,7 @@ Extending Scrapy
|
|||
topics/extensions
|
||||
topics/api
|
||||
topics/signals
|
||||
topics/scheduler
|
||||
topics/exporters
|
||||
|
||||
|
||||
|
|
@ -248,6 +249,9 @@ Extending Scrapy
|
|||
:doc:`topics/signals`
|
||||
See all available signals and how to work with them.
|
||||
|
||||
:doc:`topics/scheduler`
|
||||
Understand the scheduler component.
|
||||
|
||||
:doc:`topics/exporters`
|
||||
Quickly export your scraped items to a file (XML, CSV, etc).
|
||||
|
||||
|
|
|
|||
|
|
@ -87,8 +87,9 @@ of the system, and triggering events when certain actions occur. See the
|
|||
Scheduler
|
||||
---------
|
||||
|
||||
The Scheduler receives requests from the engine and enqueues them for feeding
|
||||
them later (also to the engine) when the engine requests them.
|
||||
The :ref:`scheduler <topics-scheduler>` receives requests from the engine and
|
||||
enqueues them for feeding them later (also to the engine) when the engine
|
||||
requests them.
|
||||
|
||||
.. _component-downloader:
|
||||
|
||||
|
|
|
|||
|
|
@ -135,6 +135,9 @@ Here are some examples to illustrate:
|
|||
|
||||
- ``s3://mybucket/scraping/feeds/%(name)s/%(time)s.json``
|
||||
|
||||
.. note:: :ref:`Spider arguments <spiderargs>` become spider attributes, hence
|
||||
they can also be used as storage URI parameters.
|
||||
|
||||
|
||||
.. _topics-feed-storage-backends:
|
||||
|
||||
|
|
|
|||
|
|
@ -119,6 +119,7 @@ Here is an example that runs multiple spiders simultaneously:
|
|||
|
||||
import scrapy
|
||||
from scrapy.crawler import CrawlerProcess
|
||||
from scrapy.utils.project import get_project_settings
|
||||
|
||||
class MySpider1(scrapy.Spider):
|
||||
# Your first spider definition
|
||||
|
|
@ -128,7 +129,8 @@ Here is an example that runs multiple spiders simultaneously:
|
|||
# Your second spider definition
|
||||
...
|
||||
|
||||
process = CrawlerProcess()
|
||||
settings = get_project_settings()
|
||||
process = CrawlerProcess(settings)
|
||||
process.crawl(MySpider1)
|
||||
process.crawl(MySpider2)
|
||||
process.start() # the script will block here until all crawling jobs are finished
|
||||
|
|
@ -141,6 +143,7 @@ Same example using :class:`~scrapy.crawler.CrawlerRunner`:
|
|||
from twisted.internet import reactor
|
||||
from scrapy.crawler import CrawlerRunner
|
||||
from scrapy.utils.log import configure_logging
|
||||
from scrapy.utils.project import get_project_settings
|
||||
|
||||
class MySpider1(scrapy.Spider):
|
||||
# Your first spider definition
|
||||
|
|
@ -151,7 +154,8 @@ Same example using :class:`~scrapy.crawler.CrawlerRunner`:
|
|||
...
|
||||
|
||||
configure_logging()
|
||||
runner = CrawlerRunner()
|
||||
settings = get_project_settings()
|
||||
runner = CrawlerRunner(settings)
|
||||
runner.crawl(MySpider1)
|
||||
runner.crawl(MySpider2)
|
||||
d = runner.join()
|
||||
|
|
@ -166,6 +170,7 @@ Same example but running the spiders sequentially by chaining the deferreds:
|
|||
from twisted.internet import reactor, defer
|
||||
from scrapy.crawler import CrawlerRunner
|
||||
from scrapy.utils.log import configure_logging
|
||||
from scrapy.utils.project import get_project_settings
|
||||
|
||||
class MySpider1(scrapy.Spider):
|
||||
# Your first spider definition
|
||||
|
|
@ -176,7 +181,8 @@ Same example but running the spiders sequentially by chaining the deferreds:
|
|||
...
|
||||
|
||||
configure_logging()
|
||||
runner = CrawlerRunner()
|
||||
settings = get_project_settings()
|
||||
runner = CrawlerRunner(settings)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def crawl():
|
||||
|
|
|
|||
|
|
@ -26,10 +26,6 @@ Request objects
|
|||
|
||||
.. autoclass:: Request
|
||||
|
||||
A :class:`Request` object represents an HTTP request, which is usually
|
||||
generated in the Spider and executed by the Downloader, and thus generating
|
||||
a :class:`Response`.
|
||||
|
||||
:param url: the URL of this request
|
||||
|
||||
If the URL is invalid, a :exc:`ValueError` exception is raised.
|
||||
|
|
@ -205,6 +201,8 @@ Request objects
|
|||
``failure.request.cb_kwargs`` in the request's errback. For more information,
|
||||
see :ref:`errback-cb_kwargs`.
|
||||
|
||||
.. autoattribute:: Request.attributes
|
||||
|
||||
.. method:: Request.copy()
|
||||
|
||||
Return a new Request which is a copy of this Request. See also:
|
||||
|
|
@ -220,6 +218,15 @@ Request objects
|
|||
|
||||
.. automethod:: from_curl
|
||||
|
||||
.. automethod:: to_dict
|
||||
|
||||
|
||||
Other functions related to requests
|
||||
-----------------------------------
|
||||
|
||||
.. autofunction:: scrapy.utils.request.request_from_dict
|
||||
|
||||
|
||||
.. _topics-request-response-ref-request-callback-arguments:
|
||||
|
||||
Passing additional data to callback functions
|
||||
|
|
@ -642,6 +649,8 @@ dealing with JSON requests.
|
|||
data into JSON format.
|
||||
:type dumps_kwargs: dict
|
||||
|
||||
.. autoattribute:: JsonRequest.attributes
|
||||
|
||||
JsonRequest usage example
|
||||
-------------------------
|
||||
|
||||
|
|
|
|||
|
|
@ -0,0 +1,34 @@
|
|||
.. _topics-scheduler:
|
||||
|
||||
=========
|
||||
Scheduler
|
||||
=========
|
||||
|
||||
.. module:: scrapy.core.scheduler
|
||||
|
||||
The scheduler component receives requests from the :ref:`engine <component-engine>`
|
||||
and stores them into persistent and/or non-persistent data structures.
|
||||
It also gets those requests and feeds them back to the engine when it
|
||||
asks for a next request to be downloaded.
|
||||
|
||||
|
||||
Overriding the default scheduler
|
||||
================================
|
||||
|
||||
You can use your own custom scheduler class by supplying its full
|
||||
Python path in the :setting:`SCHEDULER` setting.
|
||||
|
||||
|
||||
Minimal scheduler interface
|
||||
===========================
|
||||
|
||||
.. autoclass:: BaseScheduler
|
||||
:members:
|
||||
|
||||
|
||||
Default Scrapy scheduler
|
||||
========================
|
||||
|
||||
.. autoclass:: Scheduler
|
||||
:members:
|
||||
:special-members: __len__
|
||||
|
|
@ -680,12 +680,16 @@ handler (without replacement), place this in your ``settings.py``::
|
|||
|
||||
.. _http2:
|
||||
|
||||
The default HTTPS handler uses HTTP/1.1. To use HTTP/2 update
|
||||
:setting:`DOWNLOAD_HANDLERS` as follows::
|
||||
The default HTTPS handler uses HTTP/1.1. To use HTTP/2:
|
||||
|
||||
DOWNLOAD_HANDLERS = {
|
||||
'https': 'scrapy.core.downloader.handlers.http2.H2DownloadHandler',
|
||||
}
|
||||
#. Install ``Twisted[http2]>=17.9.0`` to install the packages required to
|
||||
enable HTTP/2 support in Twisted.
|
||||
|
||||
#. Update :setting:`DOWNLOAD_HANDLERS` as follows::
|
||||
|
||||
DOWNLOAD_HANDLERS = {
|
||||
'https': 'scrapy.core.downloader.handlers.http2.H2DownloadHandler',
|
||||
}
|
||||
|
||||
.. warning::
|
||||
|
||||
|
|
@ -1280,7 +1284,8 @@ SCHEDULER
|
|||
|
||||
Default: ``'scrapy.core.scheduler.Scheduler'``
|
||||
|
||||
The scheduler to use for crawling.
|
||||
The scheduler class to be used for crawling.
|
||||
See the :ref:`topics-scheduler` topic for details.
|
||||
|
||||
.. setting:: SCHEDULER_DEBUG
|
||||
|
||||
|
|
@ -1618,7 +1623,7 @@ Default: ``2083``
|
|||
Scope: ``spidermiddlewares.urllength``
|
||||
|
||||
The maximum URL length to allow for crawled URLs. For more information about
|
||||
the default value for this setting see: https://boutell.com/newfaq/misc/urllength.html
|
||||
the default value for this setting see: https://support.microsoft.com/en-us/topic/maximum-url-length-is-2-083-characters-in-internet-explorer-174e7c8a-6666-f4e0-6fd6-908b53c12246
|
||||
|
||||
.. setting:: USER_AGENT
|
||||
|
||||
|
|
@ -1641,7 +1646,6 @@ case to see how to enable and use them.
|
|||
|
||||
.. settingslist::
|
||||
|
||||
|
||||
.. _Amazon web services: https://aws.amazon.com/
|
||||
.. _breadth-first order: https://en.wikipedia.org/wiki/Breadth-first_search
|
||||
.. _depth-first order: https://en.wikipedia.org/wiki/Depth-first_search
|
||||
|
|
|
|||
|
|
@ -294,6 +294,14 @@ The above example can also be written as follows::
|
|||
def start_requests(self):
|
||||
yield scrapy.Request(f'http://www.example.com/categories/{self.category}')
|
||||
|
||||
If you are :ref:`running Scrapy from a script <run-from-script>`, you can
|
||||
specify spider arguments when calling
|
||||
:class:`CrawlerProcess.crawl <scrapy.crawler.CrawlerProcess.crawl>` or
|
||||
:class:`CrawlerRunner.crawl <scrapy.crawler.CrawlerRunner.crawl>`::
|
||||
|
||||
process = CrawlerProcess()
|
||||
process.crawl(MySpider, category="electronics")
|
||||
|
||||
Keep in mind that spider arguments are only strings.
|
||||
The spider will not do any parsing on its own.
|
||||
If you were to set the ``start_urls`` attribute from the command line,
|
||||
|
|
|
|||
|
|
@ -110,11 +110,10 @@ using the telnet console::
|
|||
Execution engine status
|
||||
|
||||
time()-engine.start_time : 8.62972998619
|
||||
engine.has_capacity() : False
|
||||
len(engine.downloader.active) : 16
|
||||
engine.scraper.is_idle() : False
|
||||
engine.spider.name : followall
|
||||
engine.spider_is_idle(engine.spider) : False
|
||||
engine.spider_is_idle() : False
|
||||
engine.slot.closing : False
|
||||
len(engine.slot.inprogress) : 16
|
||||
len(engine.slot.scheduler.dqs or []) : 0
|
||||
|
|
|
|||
1
pylintrc
1
pylintrc
|
|
@ -24,6 +24,7 @@ disable=abstract-method,
|
|||
consider-using-in,
|
||||
consider-using-set-comprehension,
|
||||
consider-using-sys-exit,
|
||||
consider-using-with,
|
||||
cyclic-import,
|
||||
dangerous-default-value,
|
||||
deprecated-method,
|
||||
|
|
|
|||
19
pytest.ini
19
pytest.ini
|
|
@ -21,20 +21,5 @@ addopts =
|
|||
markers =
|
||||
only_asyncio: marks tests as only enabled when --reactor=asyncio is passed
|
||||
only_not_asyncio: marks tests as only enabled when --reactor=asyncio is not passed
|
||||
flake8-max-line-length = 119
|
||||
flake8-ignore =
|
||||
W503
|
||||
|
||||
# Exclude files that are meant to provide top-level imports
|
||||
# E402: Module level import not at top of file
|
||||
# F401: Module imported but unused
|
||||
scrapy/__init__.py E402
|
||||
scrapy/core/downloader/handlers/http.py F401
|
||||
scrapy/http/__init__.py F401
|
||||
scrapy/linkextractors/__init__.py E402 F401
|
||||
scrapy/selector/__init__.py F401
|
||||
scrapy/spiders/__init__.py E402 F401
|
||||
|
||||
# Issues pending a review:
|
||||
scrapy/utils/url.py F403 F405
|
||||
tests/test_loader.py E741
|
||||
filterwarnings=
|
||||
ignore::DeprecationWarning:twisted.web.test.test_webclient
|
||||
|
|
|
|||
|
|
@ -98,8 +98,9 @@ class TunnelingTCP4ClientEndpoint(TCP4ClientEndpoint):
|
|||
with this endpoint comes from the pool and a CONNECT has already been issued
|
||||
for it.
|
||||
"""
|
||||
|
||||
_responseMatcher = re.compile(br'HTTP/1\.. (?P<status>\d{3})(?P<reason>.{,32})')
|
||||
_truncatedLength = 1000
|
||||
_responseAnswer = r'HTTP/1\.. (?P<status>\d{3})(?P<reason>.{,' + str(_truncatedLength) + r'})'
|
||||
_responseMatcher = re.compile(_responseAnswer.encode())
|
||||
|
||||
def __init__(self, reactor, host, port, proxyConf, contextFactory, timeout=30, bindAddress=None):
|
||||
proxyHost, proxyPort, self._proxyAuthHeader = proxyConf
|
||||
|
|
@ -144,7 +145,7 @@ class TunnelingTCP4ClientEndpoint(TCP4ClientEndpoint):
|
|||
extra = {'status': int(respm.group('status')),
|
||||
'reason': respm.group('reason').strip()}
|
||||
else:
|
||||
extra = rcvd_bytes[:32]
|
||||
extra = rcvd_bytes[:self._truncatedLength]
|
||||
self._tunnelReadyDeferred.errback(
|
||||
TunnelError('Could not open CONNECT tunnel with proxy '
|
||||
f'{self._host}:{self._port} [{extra!r}]')
|
||||
|
|
|
|||
|
|
@ -1,51 +1,62 @@
|
|||
"""
|
||||
This is the Scrapy engine which controls the Scheduler, Downloader and Spiders.
|
||||
This is the Scrapy engine which controls the Scheduler, Downloader and Spider.
|
||||
|
||||
For more information see docs/topics/architecture.rst
|
||||
|
||||
"""
|
||||
import logging
|
||||
import warnings
|
||||
from time import time
|
||||
from typing import Callable, Iterable, Iterator, Optional, Set, Union
|
||||
|
||||
from twisted.internet import defer, task
|
||||
from twisted.internet.defer import Deferred, inlineCallbacks, succeed
|
||||
from twisted.internet.task import LoopingCall
|
||||
from twisted.python.failure import Failure
|
||||
|
||||
from scrapy import signals
|
||||
from scrapy.core.scraper import Scraper
|
||||
from scrapy.exceptions import DontCloseSpider
|
||||
from scrapy.exceptions import DontCloseSpider, ScrapyDeprecationWarning
|
||||
from scrapy.http import Response, Request
|
||||
from scrapy.utils.misc import load_object
|
||||
from scrapy.utils.reactor import CallLaterOnce
|
||||
from scrapy.settings import BaseSettings
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.utils.log import logformatter_adapter, failure_to_exc_info
|
||||
from scrapy.utils.misc import create_instance, load_object
|
||||
from scrapy.utils.reactor import CallLaterOnce
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class Slot:
|
||||
|
||||
def __init__(self, start_requests, close_if_idle, nextcall, scheduler):
|
||||
self.closing = False
|
||||
self.inprogress = set() # requests in progress
|
||||
self.start_requests = iter(start_requests)
|
||||
def __init__(
|
||||
self,
|
||||
start_requests: Iterable,
|
||||
close_if_idle: bool,
|
||||
nextcall: CallLaterOnce,
|
||||
scheduler,
|
||||
) -> None:
|
||||
self.closing: Optional[Deferred] = None
|
||||
self.inprogress: Set[Request] = set()
|
||||
self.start_requests: Optional[Iterator] = iter(start_requests)
|
||||
self.close_if_idle = close_if_idle
|
||||
self.nextcall = nextcall
|
||||
self.scheduler = scheduler
|
||||
self.heartbeat = task.LoopingCall(nextcall.schedule)
|
||||
self.heartbeat = LoopingCall(nextcall.schedule)
|
||||
|
||||
def add_request(self, request):
|
||||
def add_request(self, request: Request) -> None:
|
||||
self.inprogress.add(request)
|
||||
|
||||
def remove_request(self, request):
|
||||
def remove_request(self, request: Request) -> None:
|
||||
self.inprogress.remove(request)
|
||||
self._maybe_fire_closing()
|
||||
|
||||
def close(self):
|
||||
self.closing = defer.Deferred()
|
||||
def close(self) -> Deferred:
|
||||
self.closing = Deferred()
|
||||
self._maybe_fire_closing()
|
||||
return self.closing
|
||||
|
||||
def _maybe_fire_closing(self):
|
||||
if self.closing and not self.inprogress:
|
||||
def _maybe_fire_closing(self) -> None:
|
||||
if self.closing is not None and not self.inprogress:
|
||||
if self.nextcall:
|
||||
self.nextcall.cancel()
|
||||
if self.heartbeat.running:
|
||||
|
|
@ -54,210 +65,236 @@ class Slot:
|
|||
|
||||
|
||||
class ExecutionEngine:
|
||||
|
||||
def __init__(self, crawler, spider_closed_callback):
|
||||
def __init__(self, crawler, spider_closed_callback: Callable) -> None:
|
||||
self.crawler = crawler
|
||||
self.settings = crawler.settings
|
||||
self.signals = crawler.signals
|
||||
self.logformatter = crawler.logformatter
|
||||
self.slot = None
|
||||
self.spider = None
|
||||
self.slot: Optional[Slot] = None
|
||||
self.spider: Optional[Spider] = None
|
||||
self.running = False
|
||||
self.paused = False
|
||||
self.scheduler_cls = load_object(self.settings['SCHEDULER'])
|
||||
self.scheduler_cls = self._get_scheduler_class(crawler.settings)
|
||||
downloader_cls = load_object(self.settings['DOWNLOADER'])
|
||||
self.downloader = downloader_cls(crawler)
|
||||
self.scraper = Scraper(crawler)
|
||||
self._spider_closed_callback = spider_closed_callback
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def start(self):
|
||||
"""Start the execution engine"""
|
||||
def _get_scheduler_class(self, settings: BaseSettings) -> type:
|
||||
from scrapy.core.scheduler import BaseScheduler
|
||||
scheduler_cls = load_object(settings["SCHEDULER"])
|
||||
if not issubclass(scheduler_cls, BaseScheduler):
|
||||
raise TypeError(
|
||||
f"The provided scheduler class ({settings['SCHEDULER']})"
|
||||
" does not fully implement the scheduler interface"
|
||||
)
|
||||
return scheduler_cls
|
||||
|
||||
@inlineCallbacks
|
||||
def start(self) -> Deferred:
|
||||
if self.running:
|
||||
raise RuntimeError("Engine already running")
|
||||
self.start_time = time()
|
||||
yield self.signals.send_catch_log_deferred(signal=signals.engine_started)
|
||||
self.running = True
|
||||
self._closewait = defer.Deferred()
|
||||
self._closewait = Deferred()
|
||||
yield self._closewait
|
||||
|
||||
def stop(self):
|
||||
"""Stop the execution engine gracefully"""
|
||||
def stop(self) -> Deferred:
|
||||
"""Gracefully stop the execution engine"""
|
||||
@inlineCallbacks
|
||||
def _finish_stopping_engine(_) -> Deferred:
|
||||
yield self.signals.send_catch_log_deferred(signal=signals.engine_stopped)
|
||||
self._closewait.callback(None)
|
||||
|
||||
if not self.running:
|
||||
raise RuntimeError("Engine not running")
|
||||
|
||||
self.running = False
|
||||
dfd = self._close_all_spiders()
|
||||
return dfd.addBoth(lambda _: self._finish_stopping_engine())
|
||||
dfd = self.close_spider(self.spider, reason="shutdown") if self.spider is not None else succeed(None)
|
||||
return dfd.addBoth(_finish_stopping_engine)
|
||||
|
||||
def close(self):
|
||||
"""Close the execution engine gracefully.
|
||||
|
||||
If it has already been started, stop it. In all cases, close all spiders
|
||||
and the downloader.
|
||||
def close(self) -> Deferred:
|
||||
"""
|
||||
Gracefully close the execution engine.
|
||||
If it has already been started, stop it. In all cases, close the spider and the downloader.
|
||||
"""
|
||||
if self.running:
|
||||
# Will also close spiders and downloader
|
||||
return self.stop()
|
||||
elif self.open_spiders:
|
||||
# Will also close downloader
|
||||
return self._close_all_spiders()
|
||||
else:
|
||||
return defer.succeed(self.downloader.close())
|
||||
return self.stop() # will also close spider and downloader
|
||||
if self.spider is not None:
|
||||
return self.close_spider(self.spider, reason="shutdown") # will also close downloader
|
||||
return succeed(self.downloader.close())
|
||||
|
||||
def pause(self):
|
||||
"""Pause the execution engine"""
|
||||
def pause(self) -> None:
|
||||
self.paused = True
|
||||
|
||||
def unpause(self):
|
||||
"""Resume the execution engine"""
|
||||
def unpause(self) -> None:
|
||||
self.paused = False
|
||||
|
||||
def _next_request(self, spider):
|
||||
slot = self.slot
|
||||
if not slot:
|
||||
return
|
||||
def _next_request(self) -> None:
|
||||
assert self.slot is not None # typing
|
||||
assert self.spider is not None # typing
|
||||
|
||||
if self.paused:
|
||||
return
|
||||
return None
|
||||
|
||||
while not self._needs_backout(spider):
|
||||
if not self._next_request_from_scheduler(spider):
|
||||
break
|
||||
while not self._needs_backout() and self._next_request_from_scheduler() is not None:
|
||||
pass
|
||||
|
||||
if slot.start_requests and not self._needs_backout(spider):
|
||||
if self.slot.start_requests is not None and not self._needs_backout():
|
||||
try:
|
||||
request = next(slot.start_requests)
|
||||
request = next(self.slot.start_requests)
|
||||
except StopIteration:
|
||||
slot.start_requests = None
|
||||
self.slot.start_requests = None
|
||||
except Exception:
|
||||
slot.start_requests = None
|
||||
logger.error('Error while obtaining start requests',
|
||||
exc_info=True, extra={'spider': spider})
|
||||
self.slot.start_requests = None
|
||||
logger.error('Error while obtaining start requests', exc_info=True, extra={'spider': self.spider})
|
||||
else:
|
||||
self.crawl(request, spider)
|
||||
self.crawl(request)
|
||||
|
||||
if self.spider_is_idle(spider) and slot.close_if_idle:
|
||||
self._spider_idle(spider)
|
||||
if self.spider_is_idle() and self.slot.close_if_idle:
|
||||
self._spider_idle()
|
||||
|
||||
def _needs_backout(self, spider):
|
||||
slot = self.slot
|
||||
def _needs_backout(self) -> bool:
|
||||
return (
|
||||
not self.running
|
||||
or slot.closing
|
||||
or self.slot.closing # type: ignore[union-attr]
|
||||
or self.downloader.needs_backout()
|
||||
or self.scraper.slot.needs_backout()
|
||||
or self.scraper.slot.needs_backout() # type: ignore[union-attr]
|
||||
)
|
||||
|
||||
def _next_request_from_scheduler(self, spider):
|
||||
slot = self.slot
|
||||
request = slot.scheduler.next_request()
|
||||
if not request:
|
||||
return
|
||||
d = self._download(request, spider)
|
||||
d.addBoth(self._handle_downloader_output, request, spider)
|
||||
def _next_request_from_scheduler(self) -> Optional[Deferred]:
|
||||
assert self.slot is not None # typing
|
||||
assert self.spider is not None # typing
|
||||
|
||||
request = self.slot.scheduler.next_request()
|
||||
if request is None:
|
||||
return None
|
||||
|
||||
d = self._download(request, self.spider)
|
||||
d.addBoth(self._handle_downloader_output, request)
|
||||
d.addErrback(lambda f: logger.info('Error while handling downloader output',
|
||||
exc_info=failure_to_exc_info(f),
|
||||
extra={'spider': spider}))
|
||||
d.addBoth(lambda _: slot.remove_request(request))
|
||||
extra={'spider': self.spider}))
|
||||
d.addBoth(lambda _: self.slot.remove_request(request))
|
||||
d.addErrback(lambda f: logger.info('Error while removing request from slot',
|
||||
exc_info=failure_to_exc_info(f),
|
||||
extra={'spider': spider}))
|
||||
d.addBoth(lambda _: slot.nextcall.schedule())
|
||||
extra={'spider': self.spider}))
|
||||
d.addBoth(lambda _: self.slot.nextcall.schedule())
|
||||
d.addErrback(lambda f: logger.info('Error while scheduling new request',
|
||||
exc_info=failure_to_exc_info(f),
|
||||
extra={'spider': spider}))
|
||||
extra={'spider': self.spider}))
|
||||
return d
|
||||
|
||||
def _handle_downloader_output(self, response, request, spider):
|
||||
if not isinstance(response, (Request, Response, Failure)):
|
||||
raise TypeError(
|
||||
"Incorrect type: expected Request, Response or Failure, got "
|
||||
f"{type(response)}: {response!r}"
|
||||
)
|
||||
def _handle_downloader_output(
|
||||
self, result: Union[Request, Response, Failure], request: Request
|
||||
) -> Optional[Deferred]:
|
||||
assert self.spider is not None # typing
|
||||
|
||||
if not isinstance(result, (Request, Response, Failure)):
|
||||
raise TypeError(f"Incorrect type: expected Request, Response or Failure, got {type(result)}: {result!r}")
|
||||
|
||||
# downloader middleware can return requests (for example, redirects)
|
||||
if isinstance(response, Request):
|
||||
self.crawl(response, spider)
|
||||
return
|
||||
# response is a Response or Failure
|
||||
d = self.scraper.enqueue_scrape(response, request, spider)
|
||||
d.addErrback(lambda f: logger.error('Error while enqueuing downloader output',
|
||||
exc_info=failure_to_exc_info(f),
|
||||
extra={'spider': spider}))
|
||||
if isinstance(result, Request):
|
||||
self.crawl(result)
|
||||
return None
|
||||
|
||||
d = self.scraper.enqueue_scrape(result, request, self.spider)
|
||||
d.addErrback(
|
||||
lambda f: logger.error(
|
||||
"Error while enqueuing downloader output",
|
||||
exc_info=failure_to_exc_info(f),
|
||||
extra={'spider': self.spider},
|
||||
)
|
||||
)
|
||||
return d
|
||||
|
||||
def spider_is_idle(self, spider):
|
||||
if not self.scraper.slot.is_idle():
|
||||
# scraper is not idle
|
||||
def spider_is_idle(self, spider: Optional[Spider] = None) -> bool:
|
||||
if spider is not None:
|
||||
warnings.warn(
|
||||
"Passing a 'spider' argument to ExecutionEngine.spider_is_idle is deprecated",
|
||||
category=ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
if self.slot is None:
|
||||
raise RuntimeError("Engine slot not assigned")
|
||||
if not self.scraper.slot.is_idle(): # type: ignore[union-attr]
|
||||
return False
|
||||
|
||||
if self.downloader.active:
|
||||
# downloader has pending requests
|
||||
if self.downloader.active: # downloader has pending requests
|
||||
return False
|
||||
|
||||
if self.slot.start_requests is not None:
|
||||
# not all start requests are handled
|
||||
if self.slot.start_requests is not None: # not all start requests are handled
|
||||
return False
|
||||
|
||||
if self.slot.scheduler.has_pending_requests():
|
||||
# scheduler has pending requests
|
||||
return False
|
||||
|
||||
return True
|
||||
|
||||
@property
|
||||
def open_spiders(self):
|
||||
return [self.spider] if self.spider else []
|
||||
def crawl(self, request: Request, spider: Optional[Spider] = None) -> None:
|
||||
"""Inject the request into the spider <-> downloader pipeline"""
|
||||
if spider is not None:
|
||||
warnings.warn(
|
||||
"Passing a 'spider' argument to ExecutionEngine.crawl is deprecated",
|
||||
category=ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
if spider is not self.spider:
|
||||
raise RuntimeError(f"The spider {spider.name!r} does not match the open spider")
|
||||
if self.spider is None:
|
||||
raise RuntimeError(f"No open spider to crawl: {request}")
|
||||
self._schedule_request(request, self.spider)
|
||||
self.slot.nextcall.schedule() # type: ignore[union-attr]
|
||||
|
||||
def has_capacity(self):
|
||||
"""Does the engine have capacity to handle more spiders"""
|
||||
return not bool(self.slot)
|
||||
|
||||
def crawl(self, request, spider):
|
||||
if spider not in self.open_spiders:
|
||||
raise RuntimeError(f"Spider {spider.name!r} not opened when crawling: {request}")
|
||||
self.schedule(request, spider)
|
||||
self.slot.nextcall.schedule()
|
||||
|
||||
def schedule(self, request, spider):
|
||||
def _schedule_request(self, request: Request, spider: Spider) -> None:
|
||||
self.signals.send_catch_log(signals.request_scheduled, request=request, spider=spider)
|
||||
if not self.slot.scheduler.enqueue_request(request):
|
||||
if not self.slot.scheduler.enqueue_request(request): # type: ignore[union-attr]
|
||||
self.signals.send_catch_log(signals.request_dropped, request=request, spider=spider)
|
||||
|
||||
def download(self, request, spider):
|
||||
d = self._download(request, spider)
|
||||
d.addBoth(self._downloaded, self.slot, request, spider)
|
||||
return d
|
||||
def download(self, request: Request, spider: Optional[Spider] = None) -> Deferred:
|
||||
"""Return a Deferred which fires with a Response as result, only downloader middlewares are applied"""
|
||||
if spider is None:
|
||||
spider = self.spider
|
||||
else:
|
||||
warnings.warn(
|
||||
"Passing a 'spider' argument to ExecutionEngine.download is deprecated",
|
||||
category=ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
if spider is not self.spider:
|
||||
logger.warning("The spider '%s' does not match the open spider", spider.name)
|
||||
if spider is None:
|
||||
raise RuntimeError(f"No open spider to crawl: {request}")
|
||||
return self._download(request, spider).addBoth(self._downloaded, request, spider)
|
||||
|
||||
def _downloaded(self, response, slot, request, spider):
|
||||
slot.remove_request(request)
|
||||
return self.download(response, spider) if isinstance(response, Request) else response
|
||||
def _downloaded(
|
||||
self, result: Union[Response, Request], request: Request, spider: Spider
|
||||
) -> Union[Deferred, Response]:
|
||||
assert self.slot is not None # typing
|
||||
self.slot.remove_request(request)
|
||||
return self.download(result, spider) if isinstance(result, Request) else result
|
||||
|
||||
def _download(self, request, spider):
|
||||
slot = self.slot
|
||||
slot.add_request(request)
|
||||
def _download(self, request: Request, spider: Spider) -> Deferred:
|
||||
assert self.slot is not None # typing
|
||||
|
||||
def _on_success(response):
|
||||
if not isinstance(response, (Response, Request)):
|
||||
raise TypeError(
|
||||
"Incorrect type: expected Response or Request, got "
|
||||
f"{type(response)}: {response!r}"
|
||||
)
|
||||
if isinstance(response, Response):
|
||||
if response.request is None:
|
||||
response.request = request
|
||||
logkws = self.logformatter.crawled(response.request, response, spider)
|
||||
self.slot.add_request(request)
|
||||
|
||||
def _on_success(result: Union[Response, Request]) -> Union[Response, Request]:
|
||||
if not isinstance(result, (Response, Request)):
|
||||
raise TypeError(f"Incorrect type: expected Response or Request, got {type(result)}: {result!r}")
|
||||
if isinstance(result, Response):
|
||||
if result.request is None:
|
||||
result.request = request
|
||||
logkws = self.logformatter.crawled(result.request, result, spider)
|
||||
if logkws is not None:
|
||||
logger.log(*logformatter_adapter(logkws), extra={'spider': spider})
|
||||
logger.log(*logformatter_adapter(logkws), extra={"spider": spider})
|
||||
self.signals.send_catch_log(
|
||||
signal=signals.response_received,
|
||||
response=response,
|
||||
request=response.request,
|
||||
response=result,
|
||||
request=result.request,
|
||||
spider=spider,
|
||||
)
|
||||
return response
|
||||
return result
|
||||
|
||||
def _on_complete(_):
|
||||
slot.nextcall.schedule()
|
||||
self.slot.nextcall.schedule()
|
||||
return _
|
||||
|
||||
dwld = self.downloader.fetch(request, spider)
|
||||
|
|
@ -265,58 +302,53 @@ class ExecutionEngine:
|
|||
dwld.addBoth(_on_complete)
|
||||
return dwld
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def open_spider(self, spider, start_requests=(), close_if_idle=True):
|
||||
if not self.has_capacity():
|
||||
@inlineCallbacks
|
||||
def open_spider(self, spider: Spider, start_requests: Iterable = (), close_if_idle: bool = True):
|
||||
if self.slot is not None:
|
||||
raise RuntimeError(f"No free spider slot when opening {spider.name!r}")
|
||||
logger.info("Spider opened", extra={'spider': spider})
|
||||
nextcall = CallLaterOnce(self._next_request, spider)
|
||||
scheduler = self.scheduler_cls.from_crawler(self.crawler)
|
||||
nextcall = CallLaterOnce(self._next_request)
|
||||
scheduler = create_instance(self.scheduler_cls, settings=None, crawler=self.crawler)
|
||||
start_requests = yield self.scraper.spidermw.process_start_requests(start_requests, spider)
|
||||
slot = Slot(start_requests, close_if_idle, nextcall, scheduler)
|
||||
self.slot = slot
|
||||
self.slot = Slot(start_requests, close_if_idle, nextcall, scheduler)
|
||||
self.spider = spider
|
||||
yield scheduler.open(spider)
|
||||
if hasattr(scheduler, "open"):
|
||||
yield scheduler.open(spider)
|
||||
yield self.scraper.open_spider(spider)
|
||||
self.crawler.stats.open_spider(spider)
|
||||
yield self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider)
|
||||
slot.nextcall.schedule()
|
||||
slot.heartbeat.start(5)
|
||||
self.slot.nextcall.schedule()
|
||||
self.slot.heartbeat.start(5)
|
||||
|
||||
def _spider_idle(self, spider):
|
||||
"""Called when a spider gets idle. This function is called when there
|
||||
are no remaining pages to download or schedule. It can be called
|
||||
multiple times. If some extension raises a DontCloseSpider exception
|
||||
(in the spider_idle signal handler) the spider is not closed until the
|
||||
next loop and this function is guaranteed to be called (at least) once
|
||||
again for this spider.
|
||||
def _spider_idle(self) -> None:
|
||||
"""
|
||||
res = self.signals.send_catch_log(signals.spider_idle, spider=spider, dont_log=DontCloseSpider)
|
||||
Called when a spider gets idle, i.e. when there are no remaining requests to download or schedule.
|
||||
It can be called multiple times. If a handler for the spider_idle signal raises a DontCloseSpider
|
||||
exception, the spider is not closed until the next loop and this function is guaranteed to be called
|
||||
(at least) once again.
|
||||
"""
|
||||
assert self.spider is not None # typing
|
||||
res = self.signals.send_catch_log(signals.spider_idle, spider=self.spider, dont_log=DontCloseSpider)
|
||||
if any(isinstance(x, Failure) and isinstance(x.value, DontCloseSpider) for _, x in res):
|
||||
return
|
||||
return None
|
||||
if self.spider_is_idle():
|
||||
self.close_spider(self.spider, reason='finished')
|
||||
|
||||
if self.spider_is_idle(spider):
|
||||
self.close_spider(spider, reason='finished')
|
||||
|
||||
def close_spider(self, spider, reason='cancelled'):
|
||||
def close_spider(self, spider: Spider, reason: str = "cancelled") -> Deferred:
|
||||
"""Close (cancel) spider and clear all its outstanding requests"""
|
||||
if self.slot is None:
|
||||
raise RuntimeError("Engine slot not assigned")
|
||||
|
||||
slot = self.slot
|
||||
if slot.closing:
|
||||
return slot.closing
|
||||
logger.info("Closing spider (%(reason)s)",
|
||||
{'reason': reason},
|
||||
extra={'spider': spider})
|
||||
if self.slot.closing is not None:
|
||||
return self.slot.closing
|
||||
|
||||
dfd = slot.close()
|
||||
logger.info("Closing spider (%(reason)s)", {'reason': reason}, extra={'spider': spider})
|
||||
|
||||
def log_failure(msg):
|
||||
def errback(failure):
|
||||
logger.error(
|
||||
msg,
|
||||
exc_info=failure_to_exc_info(failure),
|
||||
extra={'spider': spider}
|
||||
)
|
||||
dfd = self.slot.close()
|
||||
|
||||
def log_failure(msg: str) -> Callable:
|
||||
def errback(failure: Failure) -> None:
|
||||
logger.error(msg, exc_info=failure_to_exc_info(failure), extra={'spider': spider})
|
||||
return errback
|
||||
|
||||
dfd.addBoth(lambda _: self.downloader.close())
|
||||
|
|
@ -325,19 +357,19 @@ class ExecutionEngine:
|
|||
dfd.addBoth(lambda _: self.scraper.close_spider(spider))
|
||||
dfd.addErrback(log_failure('Scraper close failure'))
|
||||
|
||||
dfd.addBoth(lambda _: slot.scheduler.close(reason))
|
||||
dfd.addErrback(log_failure('Scheduler close failure'))
|
||||
if hasattr(self.slot.scheduler, "close"):
|
||||
dfd.addBoth(lambda _: self.slot.scheduler.close(reason))
|
||||
dfd.addErrback(log_failure("Scheduler close failure"))
|
||||
|
||||
dfd.addBoth(lambda _: self.signals.send_catch_log_deferred(
|
||||
signal=signals.spider_closed, spider=spider, reason=reason))
|
||||
signal=signals.spider_closed, spider=spider, reason=reason,
|
||||
))
|
||||
dfd.addErrback(log_failure('Error while sending spider_close signal'))
|
||||
|
||||
dfd.addBoth(lambda _: self.crawler.stats.close_spider(spider, reason=reason))
|
||||
dfd.addErrback(log_failure('Stats close failure'))
|
||||
|
||||
dfd.addBoth(lambda _: logger.info("Spider closed (%(reason)s)",
|
||||
{'reason': reason},
|
||||
extra={'spider': spider}))
|
||||
dfd.addBoth(lambda _: logger.info("Spider closed (%(reason)s)", {'reason': reason}, extra={'spider': spider}))
|
||||
|
||||
dfd.addBoth(lambda _: setattr(self, 'slot', None))
|
||||
dfd.addErrback(log_failure('Error while unassigning slot'))
|
||||
|
|
@ -349,12 +381,26 @@ class ExecutionEngine:
|
|||
|
||||
return dfd
|
||||
|
||||
def _close_all_spiders(self):
|
||||
dfds = [self.close_spider(s, reason='shutdown') for s in self.open_spiders]
|
||||
dlist = defer.DeferredList(dfds)
|
||||
return dlist
|
||||
@property
|
||||
def open_spiders(self) -> list:
|
||||
warnings.warn(
|
||||
"ExecutionEngine.open_spiders is deprecated, please use ExecutionEngine.spider instead",
|
||||
category=ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
return [self.spider] if self.spider is not None else []
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def _finish_stopping_engine(self):
|
||||
yield self.signals.send_catch_log_deferred(signal=signals.engine_stopped)
|
||||
self._closewait.callback(None)
|
||||
def has_capacity(self) -> bool:
|
||||
warnings.warn("ExecutionEngine.has_capacity is deprecated", ScrapyDeprecationWarning, stacklevel=2)
|
||||
return not bool(self.slot)
|
||||
|
||||
def schedule(self, request: Request, spider: Spider) -> None:
|
||||
warnings.warn(
|
||||
"ExecutionEngine.schedule is deprecated, please use "
|
||||
"ExecutionEngine.crawl or ExecutionEngine.download instead",
|
||||
category=ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
if self.slot is None:
|
||||
raise RuntimeError("Engine slot not assigned")
|
||||
self._schedule_request(request, spider)
|
||||
|
|
|
|||
|
|
@ -1,42 +1,179 @@
|
|||
import os
|
||||
import json
|
||||
import logging
|
||||
from os.path import join, exists
|
||||
import os
|
||||
from abc import abstractmethod
|
||||
from os.path import exists, join
|
||||
from typing import Optional, Type, TypeVar
|
||||
|
||||
from scrapy.utils.misc import load_object, create_instance
|
||||
from twisted.internet.defer import Deferred
|
||||
|
||||
from scrapy.crawler import Crawler
|
||||
from scrapy.http.request import Request
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.utils.job import job_dir
|
||||
from scrapy.utils.misc import create_instance, load_object
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class Scheduler:
|
||||
class BaseSchedulerMeta(type):
|
||||
"""
|
||||
Scrapy Scheduler. It allows to enqueue requests and then get
|
||||
a next request to download. Scheduler is also handling duplication
|
||||
filtering, via dupefilter.
|
||||
|
||||
Prioritization and queueing is not performed by the Scheduler.
|
||||
User sets ``priority`` field for each Request, and a PriorityQueue
|
||||
(defined by :setting:`SCHEDULER_PRIORITY_QUEUE`) uses these priorities
|
||||
to dequeue requests in a desired order.
|
||||
|
||||
Scheduler uses two PriorityQueue instances, configured to work in-memory
|
||||
and on-disk (optional). When on-disk queue is present, it is used by
|
||||
default, and an in-memory queue is used as a fallback for cases where
|
||||
a disk queue can't handle a request (can't serialize it).
|
||||
|
||||
:setting:`SCHEDULER_MEMORY_QUEUE` and
|
||||
:setting:`SCHEDULER_DISK_QUEUE` allow to specify lower-level queue classes
|
||||
which PriorityQueue instances would be instantiated with, to keep requests
|
||||
on disk and in memory respectively.
|
||||
|
||||
Overall, Scheduler is an object which holds several PriorityQueue instances
|
||||
(in-memory and on-disk) and implements fallback logic for them.
|
||||
Also, it handles dupefilters.
|
||||
Metaclass to check scheduler classes against the necessary interface
|
||||
"""
|
||||
def __init__(self, dupefilter, jobdir=None, dqclass=None, mqclass=None,
|
||||
logunser=False, stats=None, pqclass=None, crawler=None):
|
||||
def __instancecheck__(cls, instance):
|
||||
return cls.__subclasscheck__(type(instance))
|
||||
|
||||
def __subclasscheck__(cls, subclass):
|
||||
return (
|
||||
hasattr(subclass, "has_pending_requests") and callable(subclass.has_pending_requests)
|
||||
and hasattr(subclass, "enqueue_request") and callable(subclass.enqueue_request)
|
||||
and hasattr(subclass, "next_request") and callable(subclass.next_request)
|
||||
)
|
||||
|
||||
|
||||
class BaseScheduler(metaclass=BaseSchedulerMeta):
|
||||
"""
|
||||
The scheduler component is responsible for storing requests received from
|
||||
the engine, and feeding them back upon request (also to the engine).
|
||||
|
||||
The original sources of said requests are:
|
||||
|
||||
* Spider: ``start_requests`` method, requests created for URLs in the ``start_urls`` attribute, request callbacks
|
||||
* Spider middleware: ``process_spider_output`` and ``process_spider_exception`` methods
|
||||
* Downloader middleware: ``process_request``, ``process_response`` and ``process_exception`` methods
|
||||
|
||||
The order in which the scheduler returns its stored requests (via the ``next_request`` method)
|
||||
plays a great part in determining the order in which those requests are downloaded.
|
||||
|
||||
The methods defined in this class constitute the minimal interface that the Scrapy engine will interact with.
|
||||
"""
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler: Crawler):
|
||||
"""
|
||||
Factory method which receives the current :class:`~scrapy.crawler.Crawler` object as argument.
|
||||
"""
|
||||
return cls()
|
||||
|
||||
def open(self, spider: Spider) -> Optional[Deferred]:
|
||||
"""
|
||||
Called when the spider is opened by the engine. It receives the spider
|
||||
instance as argument and it's useful to execute initialization code.
|
||||
|
||||
:param spider: the spider object for the current crawl
|
||||
:type spider: :class:`~scrapy.spiders.Spider`
|
||||
"""
|
||||
pass
|
||||
|
||||
def close(self, reason: str) -> Optional[Deferred]:
|
||||
"""
|
||||
Called when the spider is closed by the engine. It receives the reason why the crawl
|
||||
finished as argument and it's useful to execute cleaning code.
|
||||
|
||||
:param reason: a string which describes the reason why the spider was closed
|
||||
:type reason: :class:`str`
|
||||
"""
|
||||
pass
|
||||
|
||||
@abstractmethod
|
||||
def has_pending_requests(self) -> bool:
|
||||
"""
|
||||
``True`` if the scheduler has enqueued requests, ``False`` otherwise
|
||||
"""
|
||||
raise NotImplementedError()
|
||||
|
||||
@abstractmethod
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
"""
|
||||
Process a request received by the engine.
|
||||
|
||||
Return ``True`` if the request is stored correctly, ``False`` otherwise.
|
||||
|
||||
If ``False``, the engine will fire a ``request_dropped`` signal, and
|
||||
will not make further attempts to schedule the request at a later time.
|
||||
For reference, the default Scrapy scheduler returns ``False`` when the
|
||||
request is rejected by the dupefilter.
|
||||
"""
|
||||
raise NotImplementedError()
|
||||
|
||||
@abstractmethod
|
||||
def next_request(self) -> Optional[Request]:
|
||||
"""
|
||||
Return the next :class:`~scrapy.http.Request` to be processed, or ``None``
|
||||
to indicate that there are no requests to be considered ready at the moment.
|
||||
|
||||
Returning ``None`` implies that no request from the scheduler will be sent
|
||||
to the downloader in the current reactor cycle. The engine will continue
|
||||
calling ``next_request`` until ``has_pending_requests`` is ``False``.
|
||||
"""
|
||||
raise NotImplementedError()
|
||||
|
||||
|
||||
SchedulerTV = TypeVar("SchedulerTV", bound="Scheduler")
|
||||
|
||||
|
||||
class Scheduler(BaseScheduler):
|
||||
"""
|
||||
Default Scrapy scheduler. This implementation also handles duplication
|
||||
filtering via the :setting:`dupefilter <DUPEFILTER_CLASS>`.
|
||||
|
||||
This scheduler stores requests into several priority queues (defined by the
|
||||
:setting:`SCHEDULER_PRIORITY_QUEUE` setting). In turn, said priority queues
|
||||
are backed by either memory or disk based queues (respectively defined by the
|
||||
:setting:`SCHEDULER_MEMORY_QUEUE` and :setting:`SCHEDULER_DISK_QUEUE` settings).
|
||||
|
||||
Request prioritization is almost entirely delegated to the priority queue. The only
|
||||
prioritization performed by this scheduler is using the disk-based queue if present
|
||||
(i.e. if the :setting:`JOBDIR` setting is defined) and falling back to the memory-based
|
||||
queue if a serialization error occurs. If the disk queue is not present, the memory one
|
||||
is used directly.
|
||||
|
||||
:param dupefilter: An object responsible for checking and filtering duplicate requests.
|
||||
The value for the :setting:`DUPEFILTER_CLASS` setting is used by default.
|
||||
:type dupefilter: :class:`scrapy.dupefilters.BaseDupeFilter` instance or similar:
|
||||
any class that implements the `BaseDupeFilter` interface
|
||||
|
||||
:param jobdir: The path of a directory to be used for persisting the crawl's state.
|
||||
The value for the :setting:`JOBDIR` setting is used by default.
|
||||
See :ref:`topics-jobs`.
|
||||
:type jobdir: :class:`str` or ``None``
|
||||
|
||||
:param dqclass: A class to be used as persistent request queue.
|
||||
The value for the :setting:`SCHEDULER_DISK_QUEUE` setting is used by default.
|
||||
:type dqclass: class
|
||||
|
||||
:param mqclass: A class to be used as non-persistent request queue.
|
||||
The value for the :setting:`SCHEDULER_MEMORY_QUEUE` setting is used by default.
|
||||
:type mqclass: class
|
||||
|
||||
:param logunser: A boolean that indicates whether or not unserializable requests should be logged.
|
||||
The value for the :setting:`SCHEDULER_DEBUG` setting is used by default.
|
||||
:type logunser: bool
|
||||
|
||||
:param stats: A stats collector object to record stats about the request scheduling process.
|
||||
The value for the :setting:`STATS_CLASS` setting is used by default.
|
||||
:type stats: :class:`scrapy.statscollectors.StatsCollector` instance or similar:
|
||||
any class that implements the `StatsCollector` interface
|
||||
|
||||
:param pqclass: A class to be used as priority queue for requests.
|
||||
The value for the :setting:`SCHEDULER_PRIORITY_QUEUE` setting is used by default.
|
||||
:type pqclass: class
|
||||
|
||||
:param crawler: The crawler object corresponding to the current crawl.
|
||||
:type crawler: :class:`scrapy.crawler.Crawler`
|
||||
"""
|
||||
def __init__(
|
||||
self,
|
||||
dupefilter,
|
||||
jobdir: Optional[str] = None,
|
||||
dqclass=None,
|
||||
mqclass=None,
|
||||
logunser: bool = False,
|
||||
stats=None,
|
||||
pqclass=None,
|
||||
crawler: Optional[Crawler] = None,
|
||||
):
|
||||
self.df = dupefilter
|
||||
self.dqdir = self._dqdir(jobdir)
|
||||
self.pqclass = pqclass
|
||||
|
|
@ -47,34 +184,57 @@ class Scheduler:
|
|||
self.crawler = crawler
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler):
|
||||
settings = crawler.settings
|
||||
dupefilter_cls = load_object(settings['DUPEFILTER_CLASS'])
|
||||
dupefilter = create_instance(dupefilter_cls, settings, crawler)
|
||||
pqclass = load_object(settings['SCHEDULER_PRIORITY_QUEUE'])
|
||||
dqclass = load_object(settings['SCHEDULER_DISK_QUEUE'])
|
||||
mqclass = load_object(settings['SCHEDULER_MEMORY_QUEUE'])
|
||||
logunser = settings.getbool('SCHEDULER_DEBUG')
|
||||
return cls(dupefilter, jobdir=job_dir(settings), logunser=logunser,
|
||||
stats=crawler.stats, pqclass=pqclass, dqclass=dqclass,
|
||||
mqclass=mqclass, crawler=crawler)
|
||||
def from_crawler(cls: Type[SchedulerTV], crawler) -> SchedulerTV:
|
||||
"""
|
||||
Factory method, initializes the scheduler with arguments taken from the crawl settings
|
||||
"""
|
||||
dupefilter_cls = load_object(crawler.settings['DUPEFILTER_CLASS'])
|
||||
return cls(
|
||||
dupefilter=create_instance(dupefilter_cls, crawler.settings, crawler),
|
||||
jobdir=job_dir(crawler.settings),
|
||||
dqclass=load_object(crawler.settings['SCHEDULER_DISK_QUEUE']),
|
||||
mqclass=load_object(crawler.settings['SCHEDULER_MEMORY_QUEUE']),
|
||||
logunser=crawler.settings.getbool('SCHEDULER_DEBUG'),
|
||||
stats=crawler.stats,
|
||||
pqclass=load_object(crawler.settings['SCHEDULER_PRIORITY_QUEUE']),
|
||||
crawler=crawler,
|
||||
)
|
||||
|
||||
def has_pending_requests(self):
|
||||
def has_pending_requests(self) -> bool:
|
||||
return len(self) > 0
|
||||
|
||||
def open(self, spider):
|
||||
def open(self, spider: Spider) -> Optional[Deferred]:
|
||||
"""
|
||||
(1) initialize the memory queue
|
||||
(2) initialize the disk queue if the ``jobdir`` attribute is a valid directory
|
||||
(3) return the result of the dupefilter's ``open`` method
|
||||
"""
|
||||
self.spider = spider
|
||||
self.mqs = self._mq()
|
||||
self.dqs = self._dq() if self.dqdir else None
|
||||
return self.df.open()
|
||||
|
||||
def close(self, reason):
|
||||
if self.dqs:
|
||||
def close(self, reason: str) -> Optional[Deferred]:
|
||||
"""
|
||||
(1) dump pending requests to disk if there is a disk queue
|
||||
(2) return the result of the dupefilter's ``close`` method
|
||||
"""
|
||||
if self.dqs is not None:
|
||||
state = self.dqs.close()
|
||||
assert isinstance(self.dqdir, str)
|
||||
self._write_dqs_state(self.dqdir, state)
|
||||
return self.df.close(reason)
|
||||
|
||||
def enqueue_request(self, request):
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
"""
|
||||
Unless the received request is filtered out by the Dupefilter, attempt to push
|
||||
it into the disk queue, falling back to pushing it into the memory queue.
|
||||
|
||||
Increment the appropriate stats, such as: ``scheduler/enqueued``,
|
||||
``scheduler/enqueued/disk``, ``scheduler/enqueued/memory``.
|
||||
|
||||
Return ``True`` if the request was stored successfully, ``False`` otherwise.
|
||||
"""
|
||||
if not request.dont_filter and self.df.request_seen(request):
|
||||
self.df.log(request, self.spider)
|
||||
return False
|
||||
|
|
@ -87,24 +247,35 @@ class Scheduler:
|
|||
self.stats.inc_value('scheduler/enqueued', spider=self.spider)
|
||||
return True
|
||||
|
||||
def next_request(self):
|
||||
def next_request(self) -> Optional[Request]:
|
||||
"""
|
||||
Return a :class:`~scrapy.http.Request` object from the memory queue,
|
||||
falling back to the disk queue if the memory queue is empty.
|
||||
Return ``None`` if there are no more enqueued requests.
|
||||
|
||||
Increment the appropriate stats, such as: ``scheduler/dequeued``,
|
||||
``scheduler/dequeued/disk``, ``scheduler/dequeued/memory``.
|
||||
"""
|
||||
request = self.mqs.pop()
|
||||
if request:
|
||||
if request is not None:
|
||||
self.stats.inc_value('scheduler/dequeued/memory', spider=self.spider)
|
||||
else:
|
||||
request = self._dqpop()
|
||||
if request:
|
||||
if request is not None:
|
||||
self.stats.inc_value('scheduler/dequeued/disk', spider=self.spider)
|
||||
if request:
|
||||
if request is not None:
|
||||
self.stats.inc_value('scheduler/dequeued', spider=self.spider)
|
||||
return request
|
||||
|
||||
def __len__(self):
|
||||
return len(self.dqs) + len(self.mqs) if self.dqs else len(self.mqs)
|
||||
def __len__(self) -> int:
|
||||
"""
|
||||
Return the total amount of enqueued requests
|
||||
"""
|
||||
return len(self.dqs) + len(self.mqs) if self.dqs is not None else len(self.mqs)
|
||||
|
||||
def _dqpush(self, request):
|
||||
def _dqpush(self, request: Request) -> bool:
|
||||
if self.dqs is None:
|
||||
return
|
||||
return False
|
||||
try:
|
||||
self.dqs.push(request)
|
||||
except ValueError as e: # non serializable request
|
||||
|
|
@ -115,18 +286,18 @@ class Scheduler:
|
|||
logger.warning(msg, {'request': request, 'reason': e},
|
||||
exc_info=True, extra={'spider': self.spider})
|
||||
self.logunser = False
|
||||
self.stats.inc_value('scheduler/unserializable',
|
||||
spider=self.spider)
|
||||
return
|
||||
self.stats.inc_value('scheduler/unserializable', spider=self.spider)
|
||||
return False
|
||||
else:
|
||||
return True
|
||||
|
||||
def _mqpush(self, request):
|
||||
def _mqpush(self, request: Request) -> None:
|
||||
self.mqs.push(request)
|
||||
|
||||
def _dqpop(self):
|
||||
if self.dqs:
|
||||
def _dqpop(self) -> Optional[Request]:
|
||||
if self.dqs is not None:
|
||||
return self.dqs.pop()
|
||||
return None
|
||||
|
||||
def _mq(self):
|
||||
""" Create a new priority queue instance, with in-memory storage """
|
||||
|
|
@ -150,21 +321,22 @@ class Scheduler:
|
|||
{'queuesize': len(q)}, extra={'spider': self.spider})
|
||||
return q
|
||||
|
||||
def _dqdir(self, jobdir):
|
||||
def _dqdir(self, jobdir: Optional[str]) -> Optional[str]:
|
||||
""" Return a folder name to keep disk queue state at """
|
||||
if jobdir:
|
||||
if jobdir is not None:
|
||||
dqdir = join(jobdir, 'requests.queue')
|
||||
if not exists(dqdir):
|
||||
os.makedirs(dqdir)
|
||||
return dqdir
|
||||
return None
|
||||
|
||||
def _read_dqs_state(self, dqdir):
|
||||
def _read_dqs_state(self, dqdir: str) -> list:
|
||||
path = join(dqdir, 'active.json')
|
||||
if not exists(path):
|
||||
return ()
|
||||
return []
|
||||
with open(path) as f:
|
||||
return json.load(f)
|
||||
|
||||
def _write_dqs_state(self, dqdir, state):
|
||||
def _write_dqs_state(self, dqdir: str, state: list) -> None:
|
||||
with open(join(dqdir, 'active.json'), 'w') as f:
|
||||
json.dump(state, f)
|
||||
|
|
|
|||
|
|
@ -213,7 +213,7 @@ class Scraper:
|
|||
"""
|
||||
assert self.slot is not None # typing
|
||||
if isinstance(output, Request):
|
||||
self.crawler.engine.crawl(request=output, spider=spider)
|
||||
self.crawler.engine.crawl(request=output)
|
||||
elif is_item(output):
|
||||
self.slot.itemproc_size += 1
|
||||
dfd = self.itemproc.process_item(output, spider)
|
||||
|
|
|
|||
|
|
@ -67,7 +67,7 @@ class RobotsTxtMiddleware:
|
|||
priority=self.DOWNLOAD_PRIORITY,
|
||||
meta={'dont_obey_robotstxt': True}
|
||||
)
|
||||
dfd = self.crawler.engine.download(robotsreq, spider)
|
||||
dfd = self.crawler.engine.download(robotsreq)
|
||||
dfd.addCallback(self._parse_robots, netloc, spider)
|
||||
dfd.addErrback(self._logerror, robotsreq, spider)
|
||||
dfd.addErrback(self._robots_error, netloc)
|
||||
|
|
|
|||
|
|
@ -1,35 +1,47 @@
|
|||
import os
|
||||
import logging
|
||||
import os
|
||||
from typing import Optional, Set, Type, TypeVar
|
||||
|
||||
from twisted.internet.defer import Deferred
|
||||
|
||||
from scrapy.http.request import Request
|
||||
from scrapy.settings import BaseSettings
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.utils.job import job_dir
|
||||
from scrapy.utils.request import referer_str, request_fingerprint
|
||||
|
||||
|
||||
class BaseDupeFilter:
|
||||
BaseDupeFilterTV = TypeVar("BaseDupeFilterTV", bound="BaseDupeFilter")
|
||||
|
||||
|
||||
class BaseDupeFilter:
|
||||
@classmethod
|
||||
def from_settings(cls, settings):
|
||||
def from_settings(cls: Type[BaseDupeFilterTV], settings: BaseSettings) -> BaseDupeFilterTV:
|
||||
return cls()
|
||||
|
||||
def request_seen(self, request):
|
||||
def request_seen(self, request: Request) -> bool:
|
||||
return False
|
||||
|
||||
def open(self): # can return deferred
|
||||
def open(self) -> Optional[Deferred]:
|
||||
pass
|
||||
|
||||
def close(self, reason): # can return a deferred
|
||||
def close(self, reason: str) -> Optional[Deferred]:
|
||||
pass
|
||||
|
||||
def log(self, request, spider): # log that a request has been filtered
|
||||
def log(self, request: Request, spider: Spider) -> None:
|
||||
"""Log that a request has been filtered"""
|
||||
pass
|
||||
|
||||
|
||||
RFPDupeFilterTV = TypeVar("RFPDupeFilterTV", bound="RFPDupeFilter")
|
||||
|
||||
|
||||
class RFPDupeFilter(BaseDupeFilter):
|
||||
"""Request Fingerprint duplicates filter"""
|
||||
|
||||
def __init__(self, path=None, debug=False):
|
||||
def __init__(self, path: Optional[str] = None, debug: bool = False) -> None:
|
||||
self.file = None
|
||||
self.fingerprints = set()
|
||||
self.fingerprints: Set[str] = set()
|
||||
self.logdupes = True
|
||||
self.debug = debug
|
||||
self.logger = logging.getLogger(__name__)
|
||||
|
|
@ -39,26 +51,27 @@ class RFPDupeFilter(BaseDupeFilter):
|
|||
self.fingerprints.update(x.rstrip() for x in self.file)
|
||||
|
||||
@classmethod
|
||||
def from_settings(cls, settings):
|
||||
def from_settings(cls: Type[RFPDupeFilterTV], settings: BaseSettings) -> RFPDupeFilterTV:
|
||||
debug = settings.getbool('DUPEFILTER_DEBUG')
|
||||
return cls(job_dir(settings), debug)
|
||||
|
||||
def request_seen(self, request):
|
||||
def request_seen(self, request: Request) -> bool:
|
||||
fp = self.request_fingerprint(request)
|
||||
if fp in self.fingerprints:
|
||||
return True
|
||||
self.fingerprints.add(fp)
|
||||
if self.file:
|
||||
self.file.write(fp + '\n')
|
||||
return False
|
||||
|
||||
def request_fingerprint(self, request):
|
||||
def request_fingerprint(self, request: Request) -> str:
|
||||
return request_fingerprint(request)
|
||||
|
||||
def close(self, reason):
|
||||
def close(self, reason: str) -> None:
|
||||
if self.file:
|
||||
self.file.close()
|
||||
|
||||
def log(self, request, spider):
|
||||
def log(self, request: Request, spider: Spider) -> None:
|
||||
if self.debug:
|
||||
msg = "Filtered duplicate request: %(request)s (referer: %(referer)s)"
|
||||
args = {'request': request, 'referer': referer_str(request)}
|
||||
|
|
|
|||
|
|
@ -88,10 +88,8 @@ class MemoryUsage:
|
|||
self._send_report(self.notify_mails, subj)
|
||||
self.crawler.stats.set_value('memusage/limit_notified', 1)
|
||||
|
||||
open_spiders = self.crawler.engine.open_spiders
|
||||
if open_spiders:
|
||||
for spider in open_spiders:
|
||||
self.crawler.engine.close_spider(spider, 'memusage_exceeded')
|
||||
if self.crawler.engine.spider is not None:
|
||||
self.crawler.engine.close_spider(self.crawler.engine.spider, 'memusage_exceeded')
|
||||
else:
|
||||
self.crawler.stop()
|
||||
|
||||
|
|
|
|||
|
|
@ -4,17 +4,37 @@ requests in Scrapy.
|
|||
|
||||
See documentation in docs/topics/request-response.rst
|
||||
"""
|
||||
import inspect
|
||||
from typing import Optional, Tuple
|
||||
|
||||
from w3lib.url import safe_url_string
|
||||
|
||||
import scrapy
|
||||
from scrapy.http.common import obsolete_setter
|
||||
from scrapy.http.headers import Headers
|
||||
from scrapy.utils.curl import curl_to_request_kwargs
|
||||
from scrapy.utils.python import to_bytes
|
||||
from scrapy.utils.trackref import object_ref
|
||||
from scrapy.utils.url import escape_ajax
|
||||
from scrapy.http.common import obsolete_setter
|
||||
from scrapy.utils.curl import curl_to_request_kwargs
|
||||
|
||||
|
||||
class Request(object_ref):
|
||||
"""Represents an HTTP request, which is usually generated in a Spider and
|
||||
executed by the Downloader, thus generating a :class:`Response`.
|
||||
"""
|
||||
|
||||
attributes: Tuple[str, ...] = (
|
||||
"url", "callback", "method", "headers", "body",
|
||||
"cookies", "meta", "encoding", "priority",
|
||||
"dont_filter", "errback", "flags", "cb_kwargs",
|
||||
)
|
||||
"""A tuple of :class:`str` objects containing the name of all public
|
||||
attributes of the class that are also keyword parameters of the
|
||||
``__init__`` method.
|
||||
|
||||
Currently used by :meth:`Request.replace`, :meth:`Request.to_dict` and
|
||||
:func:`~scrapy.utils.request.request_from_dict`.
|
||||
"""
|
||||
|
||||
def __init__(self, url, callback=None, method='GET', headers=None, body=None,
|
||||
cookies=None, meta=None, encoding='utf-8', priority=0,
|
||||
|
|
@ -99,11 +119,8 @@ class Request(object_ref):
|
|||
return self.replace()
|
||||
|
||||
def replace(self, *args, **kwargs):
|
||||
"""Create a new Request with the same attributes except for those
|
||||
given new values.
|
||||
"""
|
||||
for x in ['url', 'method', 'headers', 'body', 'cookies', 'meta', 'flags',
|
||||
'encoding', 'priority', 'dont_filter', 'callback', 'errback', 'cb_kwargs']:
|
||||
"""Create a new Request with the same attributes except for those given new values"""
|
||||
for x in self.attributes:
|
||||
kwargs.setdefault(x, getattr(self, x))
|
||||
cls = kwargs.pop('cls', self.__class__)
|
||||
return cls(*args, **kwargs)
|
||||
|
|
@ -136,8 +153,43 @@ class Request(object_ref):
|
|||
|
||||
To translate a cURL command into a Scrapy request,
|
||||
you may use `curl2scrapy <https://michael-shub.github.io/curl2scrapy/>`_.
|
||||
|
||||
"""
|
||||
"""
|
||||
request_kwargs = curl_to_request_kwargs(curl_command, ignore_unknown_options)
|
||||
request_kwargs.update(kwargs)
|
||||
return cls(**request_kwargs)
|
||||
|
||||
def to_dict(self, *, spider: Optional["scrapy.Spider"] = None) -> dict:
|
||||
"""Return a dictionary containing the Request's data.
|
||||
|
||||
Use :func:`~scrapy.utils.request.request_from_dict` to convert back into a :class:`~scrapy.Request` object.
|
||||
|
||||
If a spider is given, this method will try to find out the name of the spider methods used as callback
|
||||
and errback and include them in the output dict, raising an exception if they cannot be found.
|
||||
"""
|
||||
d = {
|
||||
"url": self.url, # urls are safe (safe_string_url)
|
||||
"callback": _find_method(spider, self.callback) if callable(self.callback) else self.callback,
|
||||
"errback": _find_method(spider, self.errback) if callable(self.errback) else self.errback,
|
||||
"headers": dict(self.headers),
|
||||
}
|
||||
for attr in self.attributes:
|
||||
d.setdefault(attr, getattr(self, attr))
|
||||
if type(self) is not Request:
|
||||
d["_class"] = self.__module__ + '.' + self.__class__.__name__
|
||||
return d
|
||||
|
||||
|
||||
def _find_method(obj, func):
|
||||
"""Helper function for Request.to_dict"""
|
||||
# Only instance methods contain ``__func__``
|
||||
if obj and hasattr(func, '__func__'):
|
||||
members = inspect.getmembers(obj, predicate=inspect.ismethod)
|
||||
for name, obj_func in members:
|
||||
# We need to use __func__ to access the original function object because instance
|
||||
# method objects are generated each time attribute is retrieved from instance.
|
||||
#
|
||||
# Reference: The standard type hierarchy
|
||||
# https://docs.python.org/3/reference/datamodel.html
|
||||
if obj_func.__func__ is func.__func__:
|
||||
return name
|
||||
raise ValueError(f"Function {func} is not an instance method in: {obj}")
|
||||
|
|
|
|||
|
|
@ -8,12 +8,16 @@ See documentation in docs/topics/request-response.rst
|
|||
import copy
|
||||
import json
|
||||
import warnings
|
||||
from typing import Tuple
|
||||
|
||||
from scrapy.http.request import Request
|
||||
from scrapy.utils.deprecate import create_deprecated_class
|
||||
|
||||
|
||||
class JsonRequest(Request):
|
||||
|
||||
attributes: Tuple[str, ...] = Request.attributes + ("dumps_kwargs",)
|
||||
|
||||
def __init__(self, *args, **kwargs):
|
||||
dumps_kwargs = copy.deepcopy(kwargs.pop('dumps_kwargs', {}))
|
||||
dumps_kwargs.setdefault('sort_keys', True)
|
||||
|
|
@ -36,6 +40,10 @@ class JsonRequest(Request):
|
|||
self.headers.setdefault('Content-Type', 'application/json')
|
||||
self.headers.setdefault('Accept', 'application/json, text/javascript, */*; q=0.01')
|
||||
|
||||
@property
|
||||
def dumps_kwargs(self):
|
||||
return self._dumps_kwargs
|
||||
|
||||
def replace(self, *args, **kwargs):
|
||||
body_passed = kwargs.get('body', None) is not None
|
||||
data = kwargs.pop('data', None)
|
||||
|
|
|
|||
|
|
@ -173,7 +173,7 @@ class MediaPipeline:
|
|||
errback=self.media_failed, errbackArgs=(request, info))
|
||||
else:
|
||||
self._modify_media_request(request)
|
||||
dfd = self.crawler.engine.download(request, info.spider)
|
||||
dfd = self.crawler.engine.download(request)
|
||||
dfd.addCallbacks(
|
||||
callback=self.media_downloaded, callbackArgs=(request, info), callbackKeywords={'item': item},
|
||||
errback=self.media_failed, errbackArgs=(request, info))
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@ import logging
|
|||
|
||||
from scrapy.utils.misc import create_instance
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
|
|
@ -17,8 +18,7 @@ def _path_safe(text):
|
|||
>>> _path_safe('some@symbol?').startswith('some_symbol_')
|
||||
True
|
||||
"""
|
||||
pathable_slot = "".join([c if c.isalnum() or c in '-._' else '_'
|
||||
for c in text])
|
||||
pathable_slot = "".join([c if c.isalnum() or c in '-._' else '_' for c in text])
|
||||
# as we replace some letters we can get collision for different slots
|
||||
# add we add unique part
|
||||
unique_slot = hashlib.md5(text.encode('utf8')).hexdigest()
|
||||
|
|
@ -35,6 +35,9 @@ class ScrapyPriorityQueue:
|
|||
* close()
|
||||
* __len__()
|
||||
|
||||
Optionally, the queue could provide a ``peek`` method, that should return the
|
||||
next object to be returned by ``pop``, but without removing it from the queue.
|
||||
|
||||
``__init__`` method of ScrapyPriorityQueue receives a downstream_queue_cls
|
||||
argument, which is a class used to instantiate a new (internal) queue when
|
||||
a new priority is allocated.
|
||||
|
|
@ -70,10 +73,12 @@ class ScrapyPriorityQueue:
|
|||
self.curprio = min(startprios)
|
||||
|
||||
def qfactory(self, key):
|
||||
return create_instance(self.downstream_queue_cls,
|
||||
None,
|
||||
self.crawler,
|
||||
self.key + '/' + str(key))
|
||||
return create_instance(
|
||||
self.downstream_queue_cls,
|
||||
None,
|
||||
self.crawler,
|
||||
self.key + '/' + str(key),
|
||||
)
|
||||
|
||||
def priority(self, request):
|
||||
return -request.priority
|
||||
|
|
@ -99,6 +104,18 @@ class ScrapyPriorityQueue:
|
|||
self.curprio = min(prios) if prios else None
|
||||
return m
|
||||
|
||||
def peek(self):
|
||||
"""Returns the next object to be returned by :meth:`pop`,
|
||||
but without removing it from the queue.
|
||||
|
||||
Raises :exc:`NotImplementedError` if the underlying queue class does
|
||||
not implement a ``peek`` method, which is optional for queues.
|
||||
"""
|
||||
if self.curprio is None:
|
||||
return None
|
||||
queue = self.queues[self.curprio]
|
||||
return queue.peek()
|
||||
|
||||
def close(self):
|
||||
active = []
|
||||
for p, q in self.queues.items():
|
||||
|
|
@ -116,8 +133,7 @@ class DownloaderInterface:
|
|||
self.downloader = crawler.engine.downloader
|
||||
|
||||
def stats(self, possible_slots):
|
||||
return [(self._active_downloads(slot), slot)
|
||||
for slot in possible_slots]
|
||||
return [(self._active_downloads(slot), slot) for slot in possible_slots]
|
||||
|
||||
def get_slot_key(self, request):
|
||||
return self.downloader._get_slot_key(request, None)
|
||||
|
|
@ -162,10 +178,12 @@ class DownloaderAwarePriorityQueue:
|
|||
self.pqueues[slot] = self.pqfactory(slot, startprios)
|
||||
|
||||
def pqfactory(self, slot, startprios=()):
|
||||
return ScrapyPriorityQueue(self.crawler,
|
||||
self.downstream_queue_cls,
|
||||
self.key + '/' + _path_safe(slot),
|
||||
startprios)
|
||||
return ScrapyPriorityQueue(
|
||||
self.crawler,
|
||||
self.downstream_queue_cls,
|
||||
self.key + '/' + _path_safe(slot),
|
||||
startprios,
|
||||
)
|
||||
|
||||
def pop(self):
|
||||
stats = self._downloader_interface.stats(self.pqueues)
|
||||
|
|
@ -187,9 +205,22 @@ class DownloaderAwarePriorityQueue:
|
|||
queue = self.pqueues[slot]
|
||||
queue.push(request)
|
||||
|
||||
def peek(self):
|
||||
"""Returns the next object to be returned by :meth:`pop`,
|
||||
but without removing it from the queue.
|
||||
|
||||
Raises :exc:`NotImplementedError` if the underlying queue class does
|
||||
not implement a ``peek`` method, which is optional for queues.
|
||||
"""
|
||||
stats = self._downloader_interface.stats(self.pqueues)
|
||||
if not stats:
|
||||
return None
|
||||
slot = min(stats)[1]
|
||||
queue = self.pqueues[slot]
|
||||
return queue.peek()
|
||||
|
||||
def close(self):
|
||||
active = {slot: queue.close()
|
||||
for slot, queue in self.pqueues.items()}
|
||||
active = {slot: queue.close() for slot, queue in self.pqueues.items()}
|
||||
self.pqueues.clear()
|
||||
return active
|
||||
|
||||
|
|
|
|||
|
|
@ -79,7 +79,7 @@ class Shell:
|
|||
spider = self._open_spider(request, spider)
|
||||
d = _request_deferred(request)
|
||||
d.addCallback(lambda x: (x, spider))
|
||||
self.crawler.engine.crawl(request, spider)
|
||||
self.crawler.engine.crawl(request)
|
||||
return d
|
||||
|
||||
def _open_spider(self, request, spider):
|
||||
|
|
|
|||
|
|
@ -8,7 +8,8 @@ import pickle
|
|||
|
||||
from queuelib import queue
|
||||
|
||||
from scrapy.utils.reqser import request_to_dict, request_from_dict
|
||||
from scrapy.utils.deprecate import create_deprecated_class
|
||||
from scrapy.utils.request import request_from_dict
|
||||
|
||||
|
||||
def _with_mkdir(queue_class):
|
||||
|
|
@ -19,7 +20,6 @@ def _with_mkdir(queue_class):
|
|||
dirname = os.path.dirname(path)
|
||||
if not os.path.exists(dirname):
|
||||
os.makedirs(dirname, exist_ok=True)
|
||||
|
||||
super().__init__(path, *args, **kwargs)
|
||||
|
||||
return DirectoriesCreated
|
||||
|
|
@ -38,6 +38,20 @@ def _serializable_queue(queue_class, serialize, deserialize):
|
|||
if s:
|
||||
return deserialize(s)
|
||||
|
||||
def peek(self):
|
||||
"""Returns the next object to be returned by :meth:`pop`,
|
||||
but without removing it from the queue.
|
||||
|
||||
Raises :exc:`NotImplementedError` if the underlying queue class does
|
||||
not implement a ``peek`` method, which is optional for queues.
|
||||
"""
|
||||
try:
|
||||
s = super().peek()
|
||||
except AttributeError as ex:
|
||||
raise NotImplementedError("The underlying queue class does not implement 'peek'") from ex
|
||||
if s:
|
||||
return deserialize(s)
|
||||
|
||||
return SerializableQueue
|
||||
|
||||
|
||||
|
|
@ -54,17 +68,26 @@ def _scrapy_serialization_queue(queue_class):
|
|||
return cls(crawler, key)
|
||||
|
||||
def push(self, request):
|
||||
request = request_to_dict(request, self.spider)
|
||||
request = request.to_dict(spider=self.spider)
|
||||
return super().push(request)
|
||||
|
||||
def pop(self):
|
||||
request = super().pop()
|
||||
|
||||
if not request:
|
||||
return None
|
||||
return request_from_dict(request, spider=self.spider)
|
||||
|
||||
request = request_from_dict(request, self.spider)
|
||||
return request
|
||||
def peek(self):
|
||||
"""Returns the next object to be returned by :meth:`pop`,
|
||||
but without removing it from the queue.
|
||||
|
||||
Raises :exc:`NotImplementedError` if the underlying queue class does
|
||||
not implement a ``peek`` method, which is optional for queues.
|
||||
"""
|
||||
request = super().peek()
|
||||
if not request:
|
||||
return None
|
||||
return request_from_dict(request, spider=self.spider)
|
||||
|
||||
return ScrapyRequestQueue
|
||||
|
||||
|
|
@ -76,6 +99,19 @@ def _scrapy_non_serialization_queue(queue_class):
|
|||
def from_crawler(cls, crawler, *args, **kwargs):
|
||||
return cls()
|
||||
|
||||
def peek(self):
|
||||
"""Returns the next object to be returned by :meth:`pop`,
|
||||
but without removing it from the queue.
|
||||
|
||||
Raises :exc:`NotImplementedError` if the underlying queue class does
|
||||
not implement a ``peek`` method, which is optional for queues.
|
||||
"""
|
||||
try:
|
||||
s = super().peek()
|
||||
except AttributeError as ex:
|
||||
raise NotImplementedError("The underlying queue class does not implement 'peek'") from ex
|
||||
return s
|
||||
|
||||
return ScrapyRequestQueue
|
||||
|
||||
|
||||
|
|
@ -88,38 +124,60 @@ def _pickle_serialize(obj):
|
|||
raise ValueError(str(e)) from e
|
||||
|
||||
|
||||
PickleFifoDiskQueueNonRequest = _serializable_queue(
|
||||
_PickleFifoSerializationDiskQueue = _serializable_queue(
|
||||
_with_mkdir(queue.FifoDiskQueue),
|
||||
_pickle_serialize,
|
||||
pickle.loads
|
||||
)
|
||||
PickleLifoDiskQueueNonRequest = _serializable_queue(
|
||||
_PickleLifoSerializationDiskQueue = _serializable_queue(
|
||||
_with_mkdir(queue.LifoDiskQueue),
|
||||
_pickle_serialize,
|
||||
pickle.loads
|
||||
)
|
||||
MarshalFifoDiskQueueNonRequest = _serializable_queue(
|
||||
_MarshalFifoSerializationDiskQueue = _serializable_queue(
|
||||
_with_mkdir(queue.FifoDiskQueue),
|
||||
marshal.dumps,
|
||||
marshal.loads
|
||||
)
|
||||
MarshalLifoDiskQueueNonRequest = _serializable_queue(
|
||||
_MarshalLifoSerializationDiskQueue = _serializable_queue(
|
||||
_with_mkdir(queue.LifoDiskQueue),
|
||||
marshal.dumps,
|
||||
marshal.loads
|
||||
)
|
||||
|
||||
PickleFifoDiskQueue = _scrapy_serialization_queue(
|
||||
PickleFifoDiskQueueNonRequest
|
||||
)
|
||||
PickleLifoDiskQueue = _scrapy_serialization_queue(
|
||||
PickleLifoDiskQueueNonRequest
|
||||
)
|
||||
MarshalFifoDiskQueue = _scrapy_serialization_queue(
|
||||
MarshalFifoDiskQueueNonRequest
|
||||
)
|
||||
MarshalLifoDiskQueue = _scrapy_serialization_queue(
|
||||
MarshalLifoDiskQueueNonRequest
|
||||
)
|
||||
# public queue classes
|
||||
PickleFifoDiskQueue = _scrapy_serialization_queue(_PickleFifoSerializationDiskQueue)
|
||||
PickleLifoDiskQueue = _scrapy_serialization_queue(_PickleLifoSerializationDiskQueue)
|
||||
MarshalFifoDiskQueue = _scrapy_serialization_queue(_MarshalFifoSerializationDiskQueue)
|
||||
MarshalLifoDiskQueue = _scrapy_serialization_queue(_MarshalLifoSerializationDiskQueue)
|
||||
FifoMemoryQueue = _scrapy_non_serialization_queue(queue.FifoMemoryQueue)
|
||||
LifoMemoryQueue = _scrapy_non_serialization_queue(queue.LifoMemoryQueue)
|
||||
|
||||
|
||||
# deprecated queue classes
|
||||
_subclass_warn_message = "{cls} inherits from deprecated class {old}"
|
||||
_instance_warn_message = "{cls} is deprecated"
|
||||
PickleFifoDiskQueueNonRequest = create_deprecated_class(
|
||||
name="PickleFifoDiskQueueNonRequest",
|
||||
new_class=_PickleFifoSerializationDiskQueue,
|
||||
subclass_warn_message=_subclass_warn_message,
|
||||
instance_warn_message=_instance_warn_message,
|
||||
)
|
||||
PickleLifoDiskQueueNonRequest = create_deprecated_class(
|
||||
name="PickleLifoDiskQueueNonRequest",
|
||||
new_class=_PickleLifoSerializationDiskQueue,
|
||||
subclass_warn_message=_subclass_warn_message,
|
||||
instance_warn_message=_instance_warn_message,
|
||||
)
|
||||
MarshalFifoDiskQueueNonRequest = create_deprecated_class(
|
||||
name="MarshalFifoDiskQueueNonRequest",
|
||||
new_class=_MarshalFifoSerializationDiskQueue,
|
||||
subclass_warn_message=_subclass_warn_message,
|
||||
instance_warn_message=_instance_warn_message,
|
||||
)
|
||||
MarshalLifoDiskQueueNonRequest = create_deprecated_class(
|
||||
name="MarshalLifoDiskQueueNonRequest",
|
||||
new_class=_MarshalLifoSerializationDiskQueue,
|
||||
subclass_warn_message=_subclass_warn_message,
|
||||
instance_warn_message=_instance_warn_message,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -8,11 +8,10 @@ def get_engine_status(engine):
|
|||
"""Return a report of the current engine status"""
|
||||
tests = [
|
||||
"time()-engine.start_time",
|
||||
"engine.has_capacity()",
|
||||
"len(engine.downloader.active)",
|
||||
"engine.scraper.is_idle()",
|
||||
"engine.spider.name",
|
||||
"engine.spider_is_idle(engine.spider)",
|
||||
"engine.spider_is_idle()",
|
||||
"engine.slot.closing",
|
||||
"len(engine.slot.inprogress)",
|
||||
"len(engine.slot.scheduler.dqs or [])",
|
||||
|
|
|
|||
|
|
@ -1,7 +1,10 @@
|
|||
import os
|
||||
from typing import Optional
|
||||
|
||||
from scrapy.settings import BaseSettings
|
||||
|
||||
|
||||
def job_dir(settings):
|
||||
def job_dir(settings: BaseSettings) -> Optional[str]:
|
||||
path = settings['JOBDIR']
|
||||
if path and not os.path.exists(path):
|
||||
os.makedirs(path)
|
||||
|
|
|
|||
|
|
@ -1,95 +1,22 @@
|
|||
"""
|
||||
Helper functions for serializing (and deserializing) requests.
|
||||
"""
|
||||
import inspect
|
||||
import warnings
|
||||
from typing import Optional
|
||||
|
||||
from scrapy.http import Request
|
||||
from scrapy.utils.python import to_unicode
|
||||
from scrapy.utils.misc import load_object
|
||||
import scrapy
|
||||
from scrapy.exceptions import ScrapyDeprecationWarning
|
||||
from scrapy.utils.request import request_from_dict as _from_dict
|
||||
|
||||
|
||||
def request_to_dict(request, spider=None):
|
||||
"""Convert Request object to a dict.
|
||||
|
||||
If a spider is given, it will try to find out the name of the spider method
|
||||
used in the callback and store that as the callback.
|
||||
"""
|
||||
cb = request.callback
|
||||
if callable(cb):
|
||||
cb = _find_method(spider, cb)
|
||||
eb = request.errback
|
||||
if callable(eb):
|
||||
eb = _find_method(spider, eb)
|
||||
d = {
|
||||
'url': to_unicode(request.url), # urls should be safe (safe_string_url)
|
||||
'callback': cb,
|
||||
'errback': eb,
|
||||
'method': request.method,
|
||||
'headers': dict(request.headers),
|
||||
'body': request.body,
|
||||
'cookies': request.cookies,
|
||||
'meta': request.meta,
|
||||
'_encoding': request._encoding,
|
||||
'priority': request.priority,
|
||||
'dont_filter': request.dont_filter,
|
||||
'flags': request.flags,
|
||||
'cb_kwargs': request.cb_kwargs,
|
||||
}
|
||||
if type(request) is not Request:
|
||||
d['_class'] = request.__module__ + '.' + request.__class__.__name__
|
||||
return d
|
||||
warnings.warn(
|
||||
("Module scrapy.utils.reqser is deprecated, please use request.to_dict method"
|
||||
" and/or scrapy.utils.request.request_from_dict instead"),
|
||||
category=ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
|
||||
|
||||
def request_from_dict(d, spider=None):
|
||||
"""Create Request object from a dict.
|
||||
|
||||
If a spider is given, it will try to resolve the callbacks looking at the
|
||||
spider for methods with the same name.
|
||||
"""
|
||||
cb = d['callback']
|
||||
if cb and spider:
|
||||
cb = _get_method(spider, cb)
|
||||
eb = d['errback']
|
||||
if eb and spider:
|
||||
eb = _get_method(spider, eb)
|
||||
request_cls = load_object(d['_class']) if '_class' in d else Request
|
||||
return request_cls(
|
||||
url=to_unicode(d['url']),
|
||||
callback=cb,
|
||||
errback=eb,
|
||||
method=d['method'],
|
||||
headers=d['headers'],
|
||||
body=d['body'],
|
||||
cookies=d['cookies'],
|
||||
meta=d['meta'],
|
||||
encoding=d['_encoding'],
|
||||
priority=d['priority'],
|
||||
dont_filter=d['dont_filter'],
|
||||
flags=d.get('flags'),
|
||||
cb_kwargs=d.get('cb_kwargs'),
|
||||
)
|
||||
def request_to_dict(request: "scrapy.Request", spider: Optional["scrapy.Spider"] = None) -> dict:
|
||||
return request.to_dict(spider=spider)
|
||||
|
||||
|
||||
def _find_method(obj, func):
|
||||
# Only instance methods contain ``__func__``
|
||||
if obj and hasattr(func, '__func__'):
|
||||
members = inspect.getmembers(obj, predicate=inspect.ismethod)
|
||||
for name, obj_func in members:
|
||||
# We need to use __func__ to access the original
|
||||
# function object because instance method objects
|
||||
# are generated each time attribute is retrieved from
|
||||
# instance.
|
||||
#
|
||||
# Reference: The standard type hierarchy
|
||||
# https://docs.python.org/3/reference/datamodel.html
|
||||
if obj_func.__func__ is func.__func__:
|
||||
return name
|
||||
raise ValueError(f"Function {func} is not an instance method in: {obj}")
|
||||
|
||||
|
||||
def _get_method(obj, name):
|
||||
name = str(name)
|
||||
try:
|
||||
return getattr(obj, name)
|
||||
except AttributeError:
|
||||
raise ValueError(f"Method {name!r} not found in: {obj}")
|
||||
def request_from_dict(d: dict, spider: Optional["scrapy.Spider"] = None) -> "scrapy.Request":
|
||||
return _from_dict(d, spider=spider)
|
||||
|
|
|
|||
|
|
@ -11,8 +11,9 @@ from weakref import WeakKeyDictionary
|
|||
from w3lib.http import basic_auth_header
|
||||
from w3lib.url import canonicalize_url
|
||||
|
||||
from scrapy.http import Request
|
||||
from scrapy import Request, Spider
|
||||
from scrapy.utils.httpobj import urlparse_cached
|
||||
from scrapy.utils.misc import load_object
|
||||
from scrapy.utils.python import to_bytes, to_unicode
|
||||
|
||||
|
||||
|
|
@ -24,7 +25,7 @@ def request_fingerprint(
|
|||
request: Request,
|
||||
include_headers: Optional[Iterable[Union[bytes, str]]] = None,
|
||||
keep_fragments: bool = False,
|
||||
):
|
||||
) -> str:
|
||||
"""
|
||||
Return the request fingerprint.
|
||||
|
||||
|
|
@ -106,3 +107,27 @@ def referer_str(request: Request) -> Optional[str]:
|
|||
if referrer is None:
|
||||
return referrer
|
||||
return to_unicode(referrer, errors='replace')
|
||||
|
||||
|
||||
def request_from_dict(d: dict, *, spider: Optional[Spider] = None) -> Request:
|
||||
"""Create a :class:`~scrapy.Request` object from a dict.
|
||||
|
||||
If a spider is given, it will try to resolve the callbacks looking at the
|
||||
spider for methods with the same name.
|
||||
"""
|
||||
request_cls = load_object(d["_class"]) if "_class" in d else Request
|
||||
kwargs = {key: value for key, value in d.items() if key in request_cls.attributes}
|
||||
if d.get("callback") and spider:
|
||||
kwargs["callback"] = _get_method(spider, d["callback"])
|
||||
if d.get("errback") and spider:
|
||||
kwargs["errback"] = _get_method(spider, d["errback"])
|
||||
return request_cls(**kwargs)
|
||||
|
||||
|
||||
def _get_method(obj, name):
|
||||
"""Helper function for request_from_dict"""
|
||||
name = str(name)
|
||||
try:
|
||||
return getattr(obj, name)
|
||||
except AttributeError:
|
||||
raise ValueError(f"Method {name!r} not found in: {obj}")
|
||||
|
|
|
|||
4
setup.py
4
setup.py
|
|
@ -19,7 +19,7 @@ def has_environment_marker_platform_impl_support():
|
|||
|
||||
|
||||
install_requires = [
|
||||
'Twisted[http2]>=17.9.0',
|
||||
'Twisted>=17.9.0',
|
||||
'cryptography>=2.0',
|
||||
'cssselect>=0.9.1',
|
||||
'itemloaders>=1.0.1',
|
||||
|
|
@ -31,7 +31,7 @@ install_requires = [
|
|||
'zope.interface>=4.1.3',
|
||||
'protego>=0.1.15',
|
||||
'itemadapter>=0.1.0',
|
||||
'h2>=3.0,<4.0',
|
||||
'setuptools',
|
||||
]
|
||||
extras_require = {}
|
||||
cpython_dependencies = [
|
||||
|
|
|
|||
|
|
@ -1,5 +1,5 @@
|
|||
import json
|
||||
from unittest import mock
|
||||
from unittest import mock, skipIf
|
||||
|
||||
from pytest import mark
|
||||
from testfixtures import LogCapture
|
||||
|
|
@ -7,8 +7,8 @@ from twisted.internet import defer, error, reactor
|
|||
from twisted.trial import unittest
|
||||
from twisted.web import server
|
||||
from twisted.web.error import SchemeNotSupported
|
||||
from twisted.web.http import H2_ENABLED
|
||||
|
||||
from scrapy.core.downloader.handlers.http2 import H2DownloadHandler
|
||||
from scrapy.http import Request
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.utils.misc import create_instance
|
||||
|
|
@ -21,11 +21,17 @@ from tests.test_downloader_handlers import (
|
|||
)
|
||||
|
||||
|
||||
@skipIf(not H2_ENABLED, "HTTP/2 support in Twisted is not enabled")
|
||||
class Https2TestCase(Https11TestCase):
|
||||
|
||||
scheme = 'https'
|
||||
download_handler_cls = H2DownloadHandler
|
||||
HTTP2_DATALOSS_SKIP_REASON = "Content-Length mismatch raises InvalidBodyLengthError"
|
||||
|
||||
@classmethod
|
||||
def setUpClass(cls):
|
||||
from scrapy.core.downloader.handlers.http2 import H2DownloadHandler
|
||||
cls.download_handler_cls = H2DownloadHandler
|
||||
|
||||
def test_protocol(self):
|
||||
request = Request(self.getURL("host"), method="GET")
|
||||
d = self.download_request(request, Spider("foo"))
|
||||
|
|
@ -187,9 +193,14 @@ class Https2InvalidDNSPattern(Https2TestCase):
|
|||
super(Https2InvalidDNSPattern, self).setUp()
|
||||
|
||||
|
||||
@skipIf(not H2_ENABLED, "HTTP/2 support in Twisted is not enabled")
|
||||
class Https2CustomCiphers(Https11CustomCiphers):
|
||||
scheme = 'https'
|
||||
download_handler_cls = H2DownloadHandler
|
||||
|
||||
@classmethod
|
||||
def setUpClass(cls):
|
||||
from scrapy.core.downloader.handlers.http2 import H2DownloadHandler
|
||||
cls.download_handler_cls = H2DownloadHandler
|
||||
|
||||
|
||||
class Http2MockServerTestCase(Http11MockServerTestCase):
|
||||
|
|
@ -201,6 +212,7 @@ class Http2MockServerTestCase(Http11MockServerTestCase):
|
|||
}
|
||||
|
||||
|
||||
@skipIf(not H2_ENABLED, "HTTP/2 support in Twisted is not enabled")
|
||||
class Https2ProxyTestCase(Http11ProxyTestCase):
|
||||
# only used for HTTPS tests
|
||||
keyfile = 'keys/localhost.key'
|
||||
|
|
@ -209,9 +221,13 @@ class Https2ProxyTestCase(Http11ProxyTestCase):
|
|||
scheme = 'https'
|
||||
host = u'127.0.0.1'
|
||||
|
||||
download_handler_cls = H2DownloadHandler
|
||||
expected_http_proxy_request_body = b'/'
|
||||
|
||||
@classmethod
|
||||
def setUpClass(cls):
|
||||
from scrapy.core.downloader.handlers.http2 import H2DownloadHandler
|
||||
cls.download_handler_cls = H2DownloadHandler
|
||||
|
||||
def setUp(self):
|
||||
site = server.Site(UriResource(), timeout=None)
|
||||
self.port = reactor.listenSSL(
|
||||
|
|
|
|||
|
|
@ -42,7 +42,7 @@ Disallow: /some/randome/page.html
|
|||
""".encode('utf-8')
|
||||
response = TextResponse('http://site.local/robots.txt', body=ROBOTS)
|
||||
|
||||
def return_response(request, spider):
|
||||
def return_response(request):
|
||||
deferred = Deferred()
|
||||
reactor.callFromThread(deferred.callback, response)
|
||||
return deferred
|
||||
|
|
@ -79,7 +79,7 @@ Disallow: /some/randome/page.html
|
|||
crawler.settings.set('ROBOTSTXT_OBEY', True)
|
||||
response = Response('http://site.local/robots.txt', body=b'GIF89a\xd3\x00\xfe\x00\xa2')
|
||||
|
||||
def return_response(request, spider):
|
||||
def return_response(request):
|
||||
deferred = Deferred()
|
||||
reactor.callFromThread(deferred.callback, response)
|
||||
return deferred
|
||||
|
|
@ -102,7 +102,7 @@ Disallow: /some/randome/page.html
|
|||
crawler.settings.set('ROBOTSTXT_OBEY', True)
|
||||
response = Response('http://site.local/robots.txt')
|
||||
|
||||
def return_response(request, spider):
|
||||
def return_response(request):
|
||||
deferred = Deferred()
|
||||
reactor.callFromThread(deferred.callback, response)
|
||||
return deferred
|
||||
|
|
@ -122,7 +122,7 @@ Disallow: /some/randome/page.html
|
|||
self.crawler.settings.set('ROBOTSTXT_OBEY', True)
|
||||
err = error.DNSLookupError('Robotstxt address not found')
|
||||
|
||||
def return_failure(request, spider):
|
||||
def return_failure(request):
|
||||
deferred = Deferred()
|
||||
reactor.callFromThread(deferred.errback, failure.Failure(err))
|
||||
return deferred
|
||||
|
|
@ -138,7 +138,7 @@ Disallow: /some/randome/page.html
|
|||
self.crawler.settings.set('ROBOTSTXT_OBEY', True)
|
||||
err = error.DNSLookupError('Robotstxt address not found')
|
||||
|
||||
def immediate_failure(request, spider):
|
||||
def immediate_failure(request):
|
||||
deferred = Deferred()
|
||||
deferred.errback(failure.Failure(err))
|
||||
return deferred
|
||||
|
|
@ -150,7 +150,7 @@ Disallow: /some/randome/page.html
|
|||
def test_ignore_robotstxt_request(self):
|
||||
self.crawler.settings.set('ROBOTSTXT_OBEY', True)
|
||||
|
||||
def ignore_request(request, spider):
|
||||
def ignore_request(request):
|
||||
deferred = Deferred()
|
||||
reactor.callFromThread(deferred.errback, failure.Failure(IgnoreRequest()))
|
||||
return deferred
|
||||
|
|
|
|||
|
|
@ -13,6 +13,7 @@ module with the ``runserver`` argument::
|
|||
import os
|
||||
import re
|
||||
import sys
|
||||
import warnings
|
||||
from collections import defaultdict
|
||||
from urllib.parse import urlparse
|
||||
|
||||
|
|
@ -25,6 +26,7 @@ from twisted.web import server, static, util
|
|||
|
||||
from scrapy import signals
|
||||
from scrapy.core.engine import ExecutionEngine
|
||||
from scrapy.exceptions import ScrapyDeprecationWarning
|
||||
from scrapy.http import Request
|
||||
from scrapy.item import Item, Field
|
||||
from scrapy.linkextractors import LinkExtractor
|
||||
|
|
@ -382,22 +384,104 @@ class EngineTest(unittest.TestCase):
|
|||
yield e.close()
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_close_spiders_downloader(self):
|
||||
e = ExecutionEngine(get_crawler(TestSpider), lambda _: None)
|
||||
yield e.open_spider(TestSpider(), [])
|
||||
self.assertEqual(len(e.open_spiders), 1)
|
||||
yield e.close()
|
||||
self.assertEqual(len(e.open_spiders), 0)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_close_engine_spiders_downloader(self):
|
||||
def test_start_already_running_exception(self):
|
||||
e = ExecutionEngine(get_crawler(TestSpider), lambda _: None)
|
||||
yield e.open_spider(TestSpider(), [])
|
||||
e.start()
|
||||
self.assertTrue(e.running)
|
||||
yield e.close()
|
||||
self.assertFalse(e.running)
|
||||
self.assertEqual(len(e.open_spiders), 0)
|
||||
yield self.assertFailure(e.start(), RuntimeError).addBoth(
|
||||
lambda exc: self.assertEqual(str(exc), "Engine already running")
|
||||
)
|
||||
yield e.stop()
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_close_spiders_downloader(self):
|
||||
with warnings.catch_warnings(record=True) as warning_list:
|
||||
e = ExecutionEngine(get_crawler(TestSpider), lambda _: None)
|
||||
yield e.open_spider(TestSpider(), [])
|
||||
self.assertEqual(len(e.open_spiders), 1)
|
||||
yield e.close()
|
||||
self.assertEqual(len(e.open_spiders), 0)
|
||||
self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning)
|
||||
self.assertEqual(
|
||||
str(warning_list[0].message),
|
||||
"ExecutionEngine.open_spiders is deprecated, please use ExecutionEngine.spider instead",
|
||||
)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_close_engine_spiders_downloader(self):
|
||||
with warnings.catch_warnings(record=True) as warning_list:
|
||||
e = ExecutionEngine(get_crawler(TestSpider), lambda _: None)
|
||||
yield e.open_spider(TestSpider(), [])
|
||||
e.start()
|
||||
self.assertTrue(e.running)
|
||||
yield e.close()
|
||||
self.assertFalse(e.running)
|
||||
self.assertEqual(len(e.open_spiders), 0)
|
||||
self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning)
|
||||
self.assertEqual(
|
||||
str(warning_list[0].message),
|
||||
"ExecutionEngine.open_spiders is deprecated, please use ExecutionEngine.spider instead",
|
||||
)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_crawl_deprecated_spider_arg(self):
|
||||
with warnings.catch_warnings(record=True) as warning_list:
|
||||
e = ExecutionEngine(get_crawler(TestSpider), lambda _: None)
|
||||
spider = TestSpider()
|
||||
yield e.open_spider(spider, [])
|
||||
e.start()
|
||||
e.crawl(Request("data:,"), spider)
|
||||
yield e.close()
|
||||
self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning)
|
||||
self.assertEqual(
|
||||
str(warning_list[0].message),
|
||||
"Passing a 'spider' argument to ExecutionEngine.crawl is deprecated",
|
||||
)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_download_deprecated_spider_arg(self):
|
||||
with warnings.catch_warnings(record=True) as warning_list:
|
||||
e = ExecutionEngine(get_crawler(TestSpider), lambda _: None)
|
||||
spider = TestSpider()
|
||||
yield e.open_spider(spider, [])
|
||||
e.start()
|
||||
e.download(Request("data:,"), spider)
|
||||
yield e.close()
|
||||
self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning)
|
||||
self.assertEqual(
|
||||
str(warning_list[0].message),
|
||||
"Passing a 'spider' argument to ExecutionEngine.download is deprecated",
|
||||
)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_deprecated_schedule(self):
|
||||
with warnings.catch_warnings(record=True) as warning_list:
|
||||
e = ExecutionEngine(get_crawler(TestSpider), lambda _: None)
|
||||
spider = TestSpider()
|
||||
yield e.open_spider(spider, [])
|
||||
e.start()
|
||||
e.schedule(Request("data:,"), spider)
|
||||
yield e.close()
|
||||
self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning)
|
||||
self.assertEqual(
|
||||
str(warning_list[0].message),
|
||||
"ExecutionEngine.schedule is deprecated, please use "
|
||||
"ExecutionEngine.crawl or ExecutionEngine.download instead",
|
||||
)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_deprecated_has_capacity(self):
|
||||
with warnings.catch_warnings(record=True) as warning_list:
|
||||
e = ExecutionEngine(get_crawler(TestSpider), lambda _: None)
|
||||
self.assertTrue(e.has_capacity())
|
||||
spider = TestSpider()
|
||||
yield e.open_spider(spider, [])
|
||||
self.assertFalse(e.has_capacity())
|
||||
e.start()
|
||||
yield e.close()
|
||||
self.assertTrue(e.has_capacity())
|
||||
self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning)
|
||||
self.assertEqual(str(warning_list[0].message), "ExecutionEngine.has_capacity is deprecated")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
|
|
|||
|
|
@ -6,12 +6,14 @@ import tempfile
|
|||
import unittest
|
||||
from io import BytesIO
|
||||
from datetime import datetime
|
||||
from warnings import catch_warnings, filterwarnings
|
||||
|
||||
import lxml.etree
|
||||
from itemadapter import ItemAdapter
|
||||
|
||||
from scrapy.item import Item, Field
|
||||
from scrapy.utils.python import to_unicode
|
||||
from scrapy.exceptions import ScrapyDeprecationWarning
|
||||
from scrapy.exporters import (
|
||||
BaseItemExporter, PprintItemExporter, PickleItemExporter, CsvItemExporter,
|
||||
XmlItemExporter, JsonLinesItemExporter, JsonItemExporter,
|
||||
|
|
@ -172,10 +174,12 @@ class PythonItemExporterTest(BaseItemExporterTest):
|
|||
self.assertEqual(type(exported['age'][0]['age'][0]), dict)
|
||||
|
||||
def test_export_binary(self):
|
||||
exporter = PythonItemExporter(binary=True)
|
||||
value = self.item_class(name='John\xa3', age='22')
|
||||
expected = {b'name': b'John\xc2\xa3', b'age': b'22'}
|
||||
self.assertEqual(expected, exporter.export_item(value))
|
||||
with catch_warnings():
|
||||
filterwarnings('ignore', category=ScrapyDeprecationWarning)
|
||||
exporter = PythonItemExporter(binary=True)
|
||||
value = self.item_class(name='John\xa3', age='22')
|
||||
expected = {b'name': b'John\xc2\xa3', b'age': b'22'}
|
||||
self.assertEqual(expected, exporter.export_item(value))
|
||||
|
||||
def test_nonstring_types_item(self):
|
||||
item = self._get_nonstring_types_item()
|
||||
|
|
|
|||
|
|
@ -515,7 +515,7 @@ class FromCrawlerFileFeedStorage(FileFeedStorage, FromCrawlerMixin):
|
|||
|
||||
class DummyBlockingFeedStorage(BlockingFeedStorage):
|
||||
|
||||
def __init__(self, uri):
|
||||
def __init__(self, uri, *args, feed_options=None):
|
||||
self.path = file_uri_to_path(uri)
|
||||
|
||||
def _store_in_thread(self, file):
|
||||
|
|
@ -541,7 +541,7 @@ class LogOnStoreFileStorage:
|
|||
It can be used to make sure `store` method is invoked.
|
||||
"""
|
||||
|
||||
def __init__(self, uri):
|
||||
def __init__(self, uri, feed_options=None):
|
||||
self.path = file_uri_to_path(uri)
|
||||
self.logger = getLogger()
|
||||
|
||||
|
|
|
|||
|
|
@ -5,10 +5,9 @@ import re
|
|||
import shutil
|
||||
import string
|
||||
from ipaddress import IPv4Address
|
||||
from unittest import mock
|
||||
from unittest import mock, skipIf
|
||||
from urllib.parse import urlencode
|
||||
|
||||
from h2.exceptions import InvalidBodyLengthError
|
||||
from twisted.internet import reactor
|
||||
from twisted.internet.defer import CancelledError, Deferred, DeferredList, inlineCallbacks
|
||||
from twisted.internet.endpoints import SSL4ClientEndpoint, SSL4ServerEndpoint
|
||||
|
|
@ -17,12 +16,10 @@ from twisted.internet.ssl import optionsForClientTLS, PrivateCertificate, Certif
|
|||
from twisted.python.failure import Failure
|
||||
from twisted.trial.unittest import TestCase
|
||||
from twisted.web.client import ResponseFailed, URI
|
||||
from twisted.web.http import Request as TxRequest
|
||||
from twisted.web.http import H2_ENABLED, Request as TxRequest
|
||||
from twisted.web.server import Site, NOT_DONE_YET
|
||||
from twisted.web.static import File
|
||||
|
||||
from scrapy.core.http2.protocol import H2ClientFactory, H2ClientProtocol
|
||||
from scrapy.core.http2.stream import InactiveStreamClosed, InvalidHostname
|
||||
from scrapy.http import Request, Response, JsonRequest
|
||||
from scrapy.settings import Settings
|
||||
from scrapy.spiders import Spider
|
||||
|
|
@ -173,6 +170,7 @@ def get_client_certificate(key_file, certificate_file) -> PrivateCertificate:
|
|||
return PrivateCertificate.loadPEM(pem)
|
||||
|
||||
|
||||
@skipIf(not H2_ENABLED, "HTTP/2 support in Twisted is not enabled")
|
||||
class Https2ClientProtocolTestCase(TestCase):
|
||||
scheme = 'https'
|
||||
key_file = os.path.join(os.path.dirname(__file__), 'keys', 'localhost.key')
|
||||
|
|
@ -220,6 +218,7 @@ class Https2ClientProtocolTestCase(TestCase):
|
|||
uri = URI.fromBytes(bytes(self.get_url('/'), 'utf-8'))
|
||||
|
||||
self.conn_closed_deferred = Deferred()
|
||||
from scrapy.core.http2.protocol import H2ClientFactory
|
||||
h2_client_factory = H2ClientFactory(uri, Settings(), self.conn_closed_deferred)
|
||||
client_endpoint = SSL4ClientEndpoint(reactor, self.hostname, self.port_number, client_options)
|
||||
self.client = yield client_endpoint.connect(h2_client_factory)
|
||||
|
|
@ -426,6 +425,7 @@ class Https2ClientProtocolTestCase(TestCase):
|
|||
|
||||
def assert_failure(failure: Failure):
|
||||
self.assertTrue(len(failure.value.reasons) > 0)
|
||||
from h2.exceptions import InvalidBodyLengthError
|
||||
self.assertTrue(any(
|
||||
isinstance(error, InvalidBodyLengthError)
|
||||
for error in failure.value.reasons
|
||||
|
|
@ -511,6 +511,7 @@ class Https2ClientProtocolTestCase(TestCase):
|
|||
|
||||
def assert_inactive_stream(failure):
|
||||
self.assertIsNotNone(failure.check(ResponseFailed))
|
||||
from scrapy.core.http2.stream import InactiveStreamClosed
|
||||
self.assertTrue(any(
|
||||
isinstance(e, InactiveStreamClosed)
|
||||
for e in failure.value.reasons
|
||||
|
|
@ -596,6 +597,7 @@ class Https2ClientProtocolTestCase(TestCase):
|
|||
request = Request(url)
|
||||
|
||||
def assert_invalid_hostname(failure: Failure):
|
||||
from scrapy.core.http2.stream import InvalidHostname
|
||||
self.assertIsNotNone(failure.check(InvalidHostname))
|
||||
error_msg = str(failure.value)
|
||||
self.assertIn('localhost', error_msg)
|
||||
|
|
@ -633,6 +635,7 @@ class Https2ClientProtocolTestCase(TestCase):
|
|||
|
||||
def assert_timeout_error(failure: Failure):
|
||||
for err in failure.value.reasons:
|
||||
from scrapy.core.http2.protocol import H2ClientProtocol
|
||||
if isinstance(err, TimeoutError):
|
||||
self.assertIn(f"Connection was IDLE for more than {H2ClientProtocol.IDLE_TIMEOUT}s", str(err))
|
||||
break
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
import unittest
|
||||
from unittest import mock
|
||||
from warnings import catch_warnings
|
||||
from warnings import catch_warnings, filterwarnings
|
||||
|
||||
from w3lib.encoding import resolve_encoding
|
||||
|
||||
|
|
@ -134,7 +134,9 @@ class BaseResponseTest(unittest.TestCase):
|
|||
assert isinstance(response.text, str)
|
||||
self._assert_response_encoding(response, encoding)
|
||||
self.assertEqual(response.body, body_bytes)
|
||||
self.assertEqual(response.body_as_unicode(), body_unicode)
|
||||
with catch_warnings():
|
||||
filterwarnings("ignore", category=ScrapyDeprecationWarning)
|
||||
self.assertEqual(response.body_as_unicode(), body_unicode)
|
||||
self.assertEqual(response.text, body_unicode)
|
||||
|
||||
def _assert_response_encoding(self, response, encoding):
|
||||
|
|
@ -345,8 +347,10 @@ class TextResponseTest(BaseResponseTest):
|
|||
r1 = self.response_class('http://www.example.com', body=original_string, encoding='cp1251')
|
||||
|
||||
# check body_as_unicode
|
||||
self.assertTrue(isinstance(r1.body_as_unicode(), str))
|
||||
self.assertEqual(r1.body_as_unicode(), unicode_string)
|
||||
with catch_warnings():
|
||||
filterwarnings("ignore", category=ScrapyDeprecationWarning)
|
||||
self.assertTrue(isinstance(r1.body_as_unicode(), str))
|
||||
self.assertEqual(r1.body_as_unicode(), unicode_string)
|
||||
|
||||
# check response.text
|
||||
self.assertTrue(isinstance(r1.text, str))
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
import unittest
|
||||
from unittest import mock
|
||||
from warnings import catch_warnings
|
||||
from warnings import catch_warnings, filterwarnings
|
||||
|
||||
from scrapy.exceptions import ScrapyDeprecationWarning
|
||||
from scrapy.item import ABCMeta, _BaseItem, BaseItem, DictItem, Field, Item, ItemMeta
|
||||
|
|
@ -328,16 +328,18 @@ class BaseItemTest(unittest.TestCase):
|
|||
class SubclassedItem(Item):
|
||||
pass
|
||||
|
||||
self.assertTrue(isinstance(BaseItem(), BaseItem))
|
||||
self.assertTrue(isinstance(SubclassedBaseItem(), BaseItem))
|
||||
self.assertTrue(isinstance(Item(), BaseItem))
|
||||
self.assertTrue(isinstance(SubclassedItem(), BaseItem))
|
||||
with catch_warnings():
|
||||
filterwarnings("ignore", category=ScrapyDeprecationWarning)
|
||||
self.assertTrue(isinstance(BaseItem(), BaseItem))
|
||||
self.assertTrue(isinstance(SubclassedBaseItem(), BaseItem))
|
||||
self.assertTrue(isinstance(Item(), BaseItem))
|
||||
self.assertTrue(isinstance(SubclassedItem(), BaseItem))
|
||||
|
||||
# make sure internal checks using private _BaseItem class succeed
|
||||
self.assertTrue(isinstance(BaseItem(), _BaseItem))
|
||||
self.assertTrue(isinstance(SubclassedBaseItem(), _BaseItem))
|
||||
self.assertTrue(isinstance(Item(), _BaseItem))
|
||||
self.assertTrue(isinstance(SubclassedItem(), _BaseItem))
|
||||
# make sure internal checks using private _BaseItem class succeed
|
||||
self.assertTrue(isinstance(BaseItem(), _BaseItem))
|
||||
self.assertTrue(isinstance(SubclassedBaseItem(), _BaseItem))
|
||||
self.assertTrue(isinstance(Item(), _BaseItem))
|
||||
self.assertTrue(isinstance(SubclassedItem(), _BaseItem))
|
||||
|
||||
def test_deprecation_warning(self):
|
||||
"""
|
||||
|
|
|
|||
|
|
@ -0,0 +1,144 @@
|
|||
import tempfile
|
||||
import unittest
|
||||
|
||||
import queuelib
|
||||
|
||||
from scrapy.http.request import Request
|
||||
from scrapy.pqueues import ScrapyPriorityQueue, DownloaderAwarePriorityQueue
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.squeues import FifoMemoryQueue
|
||||
from scrapy.utils.test import get_crawler
|
||||
|
||||
from tests.test_scheduler import MockDownloader, MockEngine
|
||||
|
||||
|
||||
class PriorityQueueTest(unittest.TestCase):
|
||||
def setUp(self):
|
||||
self.crawler = get_crawler(Spider)
|
||||
self.spider = self.crawler._create_spider("foo")
|
||||
|
||||
def test_queue_push_pop_one(self):
|
||||
temp_dir = tempfile.mkdtemp()
|
||||
queue = ScrapyPriorityQueue.from_crawler(self.crawler, FifoMemoryQueue, temp_dir)
|
||||
self.assertIsNone(queue.pop())
|
||||
self.assertEqual(len(queue), 0)
|
||||
req1 = Request("https://example.org/1", priority=1)
|
||||
queue.push(req1)
|
||||
self.assertEqual(len(queue), 1)
|
||||
dequeued = queue.pop()
|
||||
self.assertEqual(len(queue), 0)
|
||||
self.assertEqual(dequeued.url, req1.url)
|
||||
self.assertEqual(dequeued.priority, req1.priority)
|
||||
self.assertEqual(queue.close(), [])
|
||||
|
||||
def test_no_peek_raises(self):
|
||||
if hasattr(queuelib.queue.FifoMemoryQueue, "peek"):
|
||||
raise unittest.SkipTest("queuelib.queue.FifoMemoryQueue.peek is defined")
|
||||
temp_dir = tempfile.mkdtemp()
|
||||
queue = ScrapyPriorityQueue.from_crawler(self.crawler, FifoMemoryQueue, temp_dir)
|
||||
queue.push(Request("https://example.org"))
|
||||
with self.assertRaises(NotImplementedError, msg="The underlying queue class does not implement 'peek'"):
|
||||
queue.peek()
|
||||
queue.close()
|
||||
|
||||
def test_peek(self):
|
||||
if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"):
|
||||
raise unittest.SkipTest("queuelib.queue.FifoMemoryQueue.peek is undefined")
|
||||
temp_dir = tempfile.mkdtemp()
|
||||
queue = ScrapyPriorityQueue.from_crawler(self.crawler, FifoMemoryQueue, temp_dir)
|
||||
self.assertEqual(len(queue), 0)
|
||||
self.assertIsNone(queue.peek())
|
||||
req1 = Request("https://example.org/1")
|
||||
req2 = Request("https://example.org/2")
|
||||
req3 = Request("https://example.org/3")
|
||||
queue.push(req1)
|
||||
queue.push(req2)
|
||||
queue.push(req3)
|
||||
self.assertEqual(len(queue), 3)
|
||||
self.assertEqual(queue.peek().url, req1.url)
|
||||
self.assertEqual(queue.pop().url, req1.url)
|
||||
self.assertEqual(len(queue), 2)
|
||||
self.assertEqual(queue.peek().url, req2.url)
|
||||
self.assertEqual(queue.pop().url, req2.url)
|
||||
self.assertEqual(len(queue), 1)
|
||||
self.assertEqual(queue.peek().url, req3.url)
|
||||
self.assertEqual(queue.pop().url, req3.url)
|
||||
self.assertEqual(queue.close(), [])
|
||||
|
||||
def test_queue_push_pop_priorities(self):
|
||||
temp_dir = tempfile.mkdtemp()
|
||||
queue = ScrapyPriorityQueue.from_crawler(self.crawler, FifoMemoryQueue, temp_dir, [-1, -2, -3])
|
||||
self.assertIsNone(queue.pop())
|
||||
self.assertEqual(len(queue), 0)
|
||||
req1 = Request("https://example.org/1", priority=1)
|
||||
req2 = Request("https://example.org/2", priority=2)
|
||||
req3 = Request("https://example.org/3", priority=3)
|
||||
queue.push(req1)
|
||||
queue.push(req2)
|
||||
queue.push(req3)
|
||||
self.assertEqual(len(queue), 3)
|
||||
dequeued = queue.pop()
|
||||
self.assertEqual(len(queue), 2)
|
||||
self.assertEqual(dequeued.url, req3.url)
|
||||
self.assertEqual(dequeued.priority, req3.priority)
|
||||
self.assertEqual(queue.close(), [-1, -2])
|
||||
|
||||
|
||||
class DownloaderAwarePriorityQueueTest(unittest.TestCase):
|
||||
def setUp(self):
|
||||
crawler = get_crawler(Spider)
|
||||
crawler.engine = MockEngine(downloader=MockDownloader())
|
||||
self.queue = DownloaderAwarePriorityQueue.from_crawler(
|
||||
crawler=crawler,
|
||||
downstream_queue_cls=FifoMemoryQueue,
|
||||
key="foo/bar",
|
||||
)
|
||||
|
||||
def tearDown(self):
|
||||
self.queue.close()
|
||||
|
||||
def test_push_pop(self):
|
||||
self.assertEqual(len(self.queue), 0)
|
||||
self.assertIsNone(self.queue.pop())
|
||||
req1 = Request("http://www.example.com/1")
|
||||
req2 = Request("http://www.example.com/2")
|
||||
req3 = Request("http://www.example.com/3")
|
||||
self.queue.push(req1)
|
||||
self.queue.push(req2)
|
||||
self.queue.push(req3)
|
||||
self.assertEqual(len(self.queue), 3)
|
||||
self.assertEqual(self.queue.pop().url, req1.url)
|
||||
self.assertEqual(len(self.queue), 2)
|
||||
self.assertEqual(self.queue.pop().url, req2.url)
|
||||
self.assertEqual(len(self.queue), 1)
|
||||
self.assertEqual(self.queue.pop().url, req3.url)
|
||||
self.assertEqual(len(self.queue), 0)
|
||||
self.assertIsNone(self.queue.pop())
|
||||
|
||||
def test_no_peek_raises(self):
|
||||
if hasattr(queuelib.queue.FifoMemoryQueue, "peek"):
|
||||
raise unittest.SkipTest("queuelib.queue.FifoMemoryQueue.peek is defined")
|
||||
self.queue.push(Request("https://example.org"))
|
||||
with self.assertRaises(NotImplementedError, msg="The underlying queue class does not implement 'peek'"):
|
||||
self.queue.peek()
|
||||
|
||||
def test_peek(self):
|
||||
if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"):
|
||||
raise unittest.SkipTest("queuelib.queue.FifoMemoryQueue.peek is undefined")
|
||||
self.assertEqual(len(self.queue), 0)
|
||||
req1 = Request("https://example.org/1")
|
||||
req2 = Request("https://example.org/2")
|
||||
req3 = Request("https://example.org/3")
|
||||
self.queue.push(req1)
|
||||
self.queue.push(req2)
|
||||
self.queue.push(req3)
|
||||
self.assertEqual(len(self.queue), 3)
|
||||
self.assertEqual(self.queue.peek().url, req1.url)
|
||||
self.assertEqual(self.queue.pop().url, req1.url)
|
||||
self.assertEqual(len(self.queue), 2)
|
||||
self.assertEqual(self.queue.peek().url, req2.url)
|
||||
self.assertEqual(self.queue.pop().url, req2.url)
|
||||
self.assertEqual(len(self.queue), 1)
|
||||
self.assertEqual(self.queue.peek().url, req3.url)
|
||||
self.assertEqual(self.queue.pop().url, req3.url)
|
||||
self.assertIsNone(self.queue.peek())
|
||||
|
|
@ -1,8 +1,16 @@
|
|||
import sys
|
||||
import unittest
|
||||
import warnings
|
||||
from contextlib import suppress
|
||||
|
||||
from scrapy.http import Request, FormRequest
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.utils.reqser import request_to_dict, request_from_dict
|
||||
from scrapy import Spider, Request
|
||||
from scrapy.exceptions import ScrapyDeprecationWarning
|
||||
from scrapy.http import FormRequest, JsonRequest
|
||||
from scrapy.utils.request import request_from_dict
|
||||
|
||||
|
||||
class CustomRequest(Request):
|
||||
pass
|
||||
|
||||
|
||||
class RequestSerializationTest(unittest.TestCase):
|
||||
|
|
@ -27,7 +35,8 @@ class RequestSerializationTest(unittest.TestCase):
|
|||
priority=20,
|
||||
meta={'a': 'b'},
|
||||
cb_kwargs={'k': 'v'},
|
||||
flags=['testFlag'])
|
||||
flags=['testFlag'],
|
||||
)
|
||||
self._assert_serializes_ok(r, spider=self.spider)
|
||||
|
||||
def test_latin1_body(self):
|
||||
|
|
@ -39,7 +48,7 @@ class RequestSerializationTest(unittest.TestCase):
|
|||
self._assert_serializes_ok(r)
|
||||
|
||||
def _assert_serializes_ok(self, request, spider=None):
|
||||
d = request_to_dict(request, spider=spider)
|
||||
d = request.to_dict(spider=spider)
|
||||
request2 = request_from_dict(d, spider=spider)
|
||||
self._assert_same_request(request, request2)
|
||||
|
||||
|
|
@ -54,16 +63,21 @@ class RequestSerializationTest(unittest.TestCase):
|
|||
self.assertEqual(r1.cookies, r2.cookies)
|
||||
self.assertEqual(r1.meta, r2.meta)
|
||||
self.assertEqual(r1.cb_kwargs, r2.cb_kwargs)
|
||||
self.assertEqual(r1.encoding, r2.encoding)
|
||||
self.assertEqual(r1._encoding, r2._encoding)
|
||||
self.assertEqual(r1.priority, r2.priority)
|
||||
self.assertEqual(r1.dont_filter, r2.dont_filter)
|
||||
self.assertEqual(r1.flags, r2.flags)
|
||||
if isinstance(r1, JsonRequest):
|
||||
self.assertEqual(r1.dumps_kwargs, r2.dumps_kwargs)
|
||||
|
||||
def test_request_class(self):
|
||||
r = FormRequest("http://www.example.com")
|
||||
self._assert_serializes_ok(r, spider=self.spider)
|
||||
r = CustomRequest("http://www.example.com")
|
||||
self._assert_serializes_ok(r, spider=self.spider)
|
||||
r1 = FormRequest("http://www.example.com")
|
||||
self._assert_serializes_ok(r1, spider=self.spider)
|
||||
r2 = CustomRequest("http://www.example.com")
|
||||
self._assert_serializes_ok(r2, spider=self.spider)
|
||||
r3 = JsonRequest("http://www.example.com", dumps_kwargs={"indent": 4})
|
||||
self._assert_serializes_ok(r3, spider=self.spider)
|
||||
|
||||
def test_callback_serialization(self):
|
||||
r = Request("http://www.example.com", callback=self.spider.parse_item,
|
||||
|
|
@ -75,7 +89,7 @@ class RequestSerializationTest(unittest.TestCase):
|
|||
callback=self.spider.parse_item_reference,
|
||||
errback=self.spider.handle_error_reference)
|
||||
self._assert_serializes_ok(r, spider=self.spider)
|
||||
request_dict = request_to_dict(r, self.spider)
|
||||
request_dict = r.to_dict(spider=self.spider)
|
||||
self.assertEqual(request_dict['callback'], 'parse_item_reference')
|
||||
self.assertEqual(request_dict['errback'], 'handle_error_reference')
|
||||
|
||||
|
|
@ -84,7 +98,7 @@ class RequestSerializationTest(unittest.TestCase):
|
|||
callback=self.spider._TestSpider__parse_item_reference,
|
||||
errback=self.spider._TestSpider__handle_error_reference)
|
||||
self._assert_serializes_ok(r, spider=self.spider)
|
||||
request_dict = request_to_dict(r, self.spider)
|
||||
request_dict = r.to_dict(spider=self.spider)
|
||||
self.assertEqual(request_dict['callback'],
|
||||
'_TestSpider__parse_item_reference')
|
||||
self.assertEqual(request_dict['errback'],
|
||||
|
|
@ -110,18 +124,16 @@ class RequestSerializationTest(unittest.TestCase):
|
|||
|
||||
def test_unserializable_callback1(self):
|
||||
r = Request("http://www.example.com", callback=lambda x: x)
|
||||
self.assertRaises(ValueError, request_to_dict, r)
|
||||
self.assertRaises(ValueError, request_to_dict, r, spider=self.spider)
|
||||
self.assertRaises(ValueError, r.to_dict, spider=self.spider)
|
||||
|
||||
def test_unserializable_callback2(self):
|
||||
r = Request("http://www.example.com", callback=self.spider.parse_item)
|
||||
self.assertRaises(ValueError, request_to_dict, r)
|
||||
self.assertRaises(ValueError, r.to_dict, spider=None)
|
||||
|
||||
def test_unserializable_callback3(self):
|
||||
"""Parser method is removed or replaced dynamically."""
|
||||
|
||||
class MySpider(Spider):
|
||||
|
||||
name = 'my_spider'
|
||||
|
||||
def parse(self, response):
|
||||
|
|
@ -130,7 +142,35 @@ class RequestSerializationTest(unittest.TestCase):
|
|||
spider = MySpider()
|
||||
r = Request("http://www.example.com", callback=spider.parse)
|
||||
setattr(spider, 'parse', None)
|
||||
self.assertRaises(ValueError, request_to_dict, r, spider=spider)
|
||||
self.assertRaises(ValueError, r.to_dict, spider=spider)
|
||||
|
||||
def test_callback_not_available(self):
|
||||
"""Callback method is not available in the spider passed to from_dict"""
|
||||
spider = TestSpiderDelegation()
|
||||
r = Request("http://www.example.com", callback=spider.delegated_callback)
|
||||
d = r.to_dict(spider=spider)
|
||||
self.assertRaises(ValueError, request_from_dict, d, spider=Spider("foo"))
|
||||
|
||||
|
||||
class DeprecatedMethodsRequestSerializationTest(RequestSerializationTest):
|
||||
def _assert_serializes_ok(self, request, spider=None):
|
||||
with warnings.catch_warnings(record=True) as caught:
|
||||
warnings.simplefilter("always")
|
||||
with suppress(KeyError):
|
||||
del sys.modules["scrapy.utils.reqser"] # delete module to reset the deprecation warning
|
||||
|
||||
from scrapy.utils.reqser import request_from_dict as _from_dict, request_to_dict as _to_dict
|
||||
|
||||
request_copy = _from_dict(_to_dict(request, spider), spider)
|
||||
self._assert_same_request(request, request_copy)
|
||||
|
||||
self.assertEqual(len(caught), 1)
|
||||
self.assertTrue(issubclass(caught[0].category, ScrapyDeprecationWarning))
|
||||
self.assertEqual(
|
||||
"Module scrapy.utils.reqser is deprecated, please use request.to_dict method"
|
||||
" and/or scrapy.utils.request.request_from_dict instead",
|
||||
str(caught[0].message),
|
||||
)
|
||||
|
||||
|
||||
class TestSpiderMixin:
|
||||
|
|
@ -177,7 +217,3 @@ class TestSpider(Spider, TestSpiderMixin):
|
|||
|
||||
def __parse_item_private(self, response):
|
||||
pass
|
||||
|
||||
|
||||
class CustomRequest(Request):
|
||||
pass
|
||||
|
|
@ -0,0 +1,159 @@
|
|||
from typing import Dict, Optional
|
||||
from unittest import TestCase
|
||||
from urllib.parse import urljoin, urlparse
|
||||
|
||||
from testfixtures import LogCapture
|
||||
from twisted.internet import defer
|
||||
from twisted.trial.unittest import TestCase as TwistedTestCase
|
||||
|
||||
from scrapy.core.scheduler import BaseScheduler
|
||||
from scrapy.crawler import CrawlerRunner
|
||||
from scrapy.http import Request
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.utils.request import request_fingerprint
|
||||
|
||||
from tests.mockserver import MockServer
|
||||
|
||||
|
||||
PATHS = ["/a", "/b", "/c"]
|
||||
URLS = [urljoin("https://example.org", p) for p in PATHS]
|
||||
|
||||
|
||||
class MinimalScheduler:
|
||||
def __init__(self) -> None:
|
||||
self.requests: Dict[str, Request] = {}
|
||||
|
||||
def has_pending_requests(self) -> bool:
|
||||
return bool(self.requests)
|
||||
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
fp = request_fingerprint(request)
|
||||
if fp not in self.requests:
|
||||
self.requests[fp] = request
|
||||
return True
|
||||
return False
|
||||
|
||||
def next_request(self) -> Optional[Request]:
|
||||
if self.has_pending_requests():
|
||||
fp, request = self.requests.popitem()
|
||||
return request
|
||||
return None
|
||||
|
||||
|
||||
class SimpleScheduler(MinimalScheduler):
|
||||
def open(self, spider: Spider) -> defer.Deferred:
|
||||
return defer.succeed("open")
|
||||
|
||||
def close(self, reason: str) -> defer.Deferred:
|
||||
return defer.succeed("close")
|
||||
|
||||
def __len__(self) -> int:
|
||||
return len(self.requests)
|
||||
|
||||
|
||||
class TestSpider(Spider):
|
||||
name = "test"
|
||||
|
||||
def __init__(self, mockserver, *args, **kwargs):
|
||||
super().__init__(*args, **kwargs)
|
||||
self.start_urls = map(mockserver.url, PATHS)
|
||||
|
||||
def parse(self, response):
|
||||
return {"path": urlparse(response.url).path}
|
||||
|
||||
|
||||
class InterfaceCheckMixin:
|
||||
def test_scheduler_class(self):
|
||||
self.assertTrue(isinstance(self.scheduler, BaseScheduler))
|
||||
self.assertTrue(issubclass(self.scheduler.__class__, BaseScheduler))
|
||||
|
||||
|
||||
class BaseSchedulerTest(TestCase, InterfaceCheckMixin):
|
||||
def setUp(self):
|
||||
self.scheduler = BaseScheduler()
|
||||
|
||||
def test_methods(self):
|
||||
self.assertIsNone(self.scheduler.open(Spider("foo")))
|
||||
self.assertIsNone(self.scheduler.close("finished"))
|
||||
self.assertRaises(NotImplementedError, self.scheduler.has_pending_requests)
|
||||
self.assertRaises(NotImplementedError, self.scheduler.enqueue_request, Request("https://example.org"))
|
||||
self.assertRaises(NotImplementedError, self.scheduler.next_request)
|
||||
|
||||
|
||||
class MinimalSchedulerTest(TestCase, InterfaceCheckMixin):
|
||||
def setUp(self):
|
||||
self.scheduler = MinimalScheduler()
|
||||
|
||||
def test_open_close(self):
|
||||
with self.assertRaises(AttributeError):
|
||||
self.scheduler.open(Spider("foo"))
|
||||
with self.assertRaises(AttributeError):
|
||||
self.scheduler.close("finished")
|
||||
|
||||
def test_len(self):
|
||||
with self.assertRaises(AttributeError):
|
||||
self.scheduler.__len__()
|
||||
with self.assertRaises(TypeError):
|
||||
len(self.scheduler)
|
||||
|
||||
def test_enqueue_dequeue(self):
|
||||
self.assertFalse(self.scheduler.has_pending_requests())
|
||||
for url in URLS:
|
||||
self.assertTrue(self.scheduler.enqueue_request(Request(url)))
|
||||
self.assertFalse(self.scheduler.enqueue_request(Request(url)))
|
||||
self.assertTrue(self.scheduler.has_pending_requests)
|
||||
|
||||
dequeued = []
|
||||
while self.scheduler.has_pending_requests():
|
||||
request = self.scheduler.next_request()
|
||||
dequeued.append(request.url)
|
||||
self.assertEqual(set(dequeued), set(URLS))
|
||||
self.assertFalse(self.scheduler.has_pending_requests())
|
||||
|
||||
|
||||
class SimpleSchedulerTest(TwistedTestCase, InterfaceCheckMixin):
|
||||
def setUp(self):
|
||||
self.scheduler = SimpleScheduler()
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_enqueue_dequeue(self):
|
||||
open_result = yield self.scheduler.open(Spider("foo"))
|
||||
self.assertEqual(open_result, "open")
|
||||
self.assertFalse(self.scheduler.has_pending_requests())
|
||||
|
||||
for url in URLS:
|
||||
self.assertTrue(self.scheduler.enqueue_request(Request(url)))
|
||||
self.assertFalse(self.scheduler.enqueue_request(Request(url)))
|
||||
|
||||
self.assertTrue(self.scheduler.has_pending_requests())
|
||||
self.assertEqual(len(self.scheduler), len(URLS))
|
||||
|
||||
dequeued = []
|
||||
while self.scheduler.has_pending_requests():
|
||||
request = self.scheduler.next_request()
|
||||
dequeued.append(request.url)
|
||||
self.assertEqual(set(dequeued), set(URLS))
|
||||
|
||||
self.assertFalse(self.scheduler.has_pending_requests())
|
||||
self.assertEqual(len(self.scheduler), 0)
|
||||
|
||||
close_result = yield self.scheduler.close("")
|
||||
self.assertEqual(close_result, "close")
|
||||
|
||||
|
||||
class MinimalSchedulerCrawlTest(TwistedTestCase):
|
||||
scheduler_cls = MinimalScheduler
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_crawl(self):
|
||||
with MockServer() as mockserver:
|
||||
settings = {"SCHEDULER": self.scheduler_cls}
|
||||
with LogCapture() as log:
|
||||
yield CrawlerRunner(settings).crawl(TestSpider, mockserver)
|
||||
for path in PATHS:
|
||||
self.assertIn(f"{{'path': '{path}'}}", str(log))
|
||||
self.assertIn(f"'item_scraped_count': {len(PATHS)}", str(log))
|
||||
|
||||
|
||||
class SimpleSchedulerCrawlTest(MinimalSchedulerCrawlTest):
|
||||
scheduler_cls = SimpleScheduler
|
||||
|
|
@ -3,10 +3,10 @@ import sys
|
|||
|
||||
from queuelib.tests import test_queue as t
|
||||
from scrapy.squeues import (
|
||||
MarshalFifoDiskQueueNonRequest as MarshalFifoDiskQueue,
|
||||
MarshalLifoDiskQueueNonRequest as MarshalLifoDiskQueue,
|
||||
PickleFifoDiskQueueNonRequest as PickleFifoDiskQueue,
|
||||
PickleLifoDiskQueueNonRequest as PickleLifoDiskQueue
|
||||
_MarshalFifoSerializationDiskQueue,
|
||||
_MarshalLifoSerializationDiskQueue,
|
||||
_PickleFifoSerializationDiskQueue,
|
||||
_PickleLifoSerializationDiskQueue,
|
||||
)
|
||||
from scrapy.item import Item, Field
|
||||
from scrapy.http import Request
|
||||
|
|
@ -53,7 +53,7 @@ class MarshalFifoDiskQueueTest(t.FifoDiskQueueTest, FifoDiskQueueTestMixin):
|
|||
chunksize = 100000
|
||||
|
||||
def queue(self):
|
||||
return MarshalFifoDiskQueue(self.qpath, chunksize=self.chunksize)
|
||||
return _MarshalFifoSerializationDiskQueue(self.qpath, chunksize=self.chunksize)
|
||||
|
||||
|
||||
class ChunkSize1MarshalFifoDiskQueueTest(MarshalFifoDiskQueueTest):
|
||||
|
|
@ -77,7 +77,7 @@ class PickleFifoDiskQueueTest(t.FifoDiskQueueTest, FifoDiskQueueTestMixin):
|
|||
chunksize = 100000
|
||||
|
||||
def queue(self):
|
||||
return PickleFifoDiskQueue(self.qpath, chunksize=self.chunksize)
|
||||
return _PickleFifoSerializationDiskQueue(self.qpath, chunksize=self.chunksize)
|
||||
|
||||
def test_serialize_item(self):
|
||||
q = self.queue()
|
||||
|
|
@ -155,13 +155,13 @@ class LifoDiskQueueTestMixin:
|
|||
class MarshalLifoDiskQueueTest(t.LifoDiskQueueTest, LifoDiskQueueTestMixin):
|
||||
|
||||
def queue(self):
|
||||
return MarshalLifoDiskQueue(self.qpath)
|
||||
return _MarshalLifoSerializationDiskQueue(self.qpath)
|
||||
|
||||
|
||||
class PickleLifoDiskQueueTest(t.LifoDiskQueueTest, LifoDiskQueueTestMixin):
|
||||
|
||||
def queue(self):
|
||||
return PickleLifoDiskQueue(self.qpath)
|
||||
return _PickleLifoSerializationDiskQueue(self.qpath)
|
||||
|
||||
def test_serialize_item(self):
|
||||
q = self.queue()
|
||||
|
|
|
|||
|
|
@ -0,0 +1,214 @@
|
|||
import shutil
|
||||
import tempfile
|
||||
import unittest
|
||||
|
||||
import queuelib
|
||||
|
||||
from scrapy.squeues import (
|
||||
PickleFifoDiskQueue,
|
||||
PickleLifoDiskQueue,
|
||||
MarshalFifoDiskQueue,
|
||||
MarshalLifoDiskQueue,
|
||||
FifoMemoryQueue,
|
||||
LifoMemoryQueue,
|
||||
)
|
||||
from scrapy.http import Request
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.utils.test import get_crawler
|
||||
|
||||
|
||||
"""
|
||||
Queues that handle requests
|
||||
"""
|
||||
|
||||
|
||||
class BaseQueueTestCase(unittest.TestCase):
|
||||
def setUp(self):
|
||||
self.tmpdir = tempfile.mkdtemp(prefix="scrapy-queue-tests-")
|
||||
self.qpath = self.tempfilename()
|
||||
self.qdir = self.mkdtemp()
|
||||
self.crawler = get_crawler(Spider)
|
||||
|
||||
def tearDown(self):
|
||||
shutil.rmtree(self.tmpdir)
|
||||
|
||||
def tempfilename(self):
|
||||
with tempfile.NamedTemporaryFile(dir=self.tmpdir) as nf:
|
||||
return nf.name
|
||||
|
||||
def mkdtemp(self):
|
||||
return tempfile.mkdtemp(dir=self.tmpdir)
|
||||
|
||||
|
||||
class RequestQueueTestMixin:
|
||||
def queue(self):
|
||||
raise NotImplementedError()
|
||||
|
||||
def test_one_element_with_peek(self):
|
||||
if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"):
|
||||
raise unittest.SkipTest("The queuelib queues do not define peek")
|
||||
q = self.queue()
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.peek())
|
||||
self.assertIsNone(q.pop())
|
||||
req = Request("http://www.example.com")
|
||||
q.push(req)
|
||||
self.assertEqual(len(q), 1)
|
||||
self.assertEqual(q.peek().url, req.url)
|
||||
self.assertEqual(q.pop().url, req.url)
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.peek())
|
||||
self.assertIsNone(q.pop())
|
||||
q.close()
|
||||
|
||||
def test_one_element_without_peek(self):
|
||||
if hasattr(queuelib.queue.FifoMemoryQueue, "peek"):
|
||||
raise unittest.SkipTest("The queuelib queues define peek")
|
||||
q = self.queue()
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.pop())
|
||||
req = Request("http://www.example.com")
|
||||
q.push(req)
|
||||
self.assertEqual(len(q), 1)
|
||||
with self.assertRaises(NotImplementedError, msg="The underlying queue class does not implement 'peek'"):
|
||||
q.peek()
|
||||
self.assertEqual(q.pop().url, req.url)
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.pop())
|
||||
q.close()
|
||||
|
||||
|
||||
class FifoQueueMixin(RequestQueueTestMixin):
|
||||
def test_fifo_with_peek(self):
|
||||
if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"):
|
||||
raise unittest.SkipTest("The queuelib queues do not define peek")
|
||||
q = self.queue()
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.peek())
|
||||
self.assertIsNone(q.pop())
|
||||
req1 = Request("http://www.example.com/1")
|
||||
req2 = Request("http://www.example.com/2")
|
||||
req3 = Request("http://www.example.com/3")
|
||||
q.push(req1)
|
||||
q.push(req2)
|
||||
q.push(req3)
|
||||
self.assertEqual(len(q), 3)
|
||||
self.assertEqual(q.peek().url, req1.url)
|
||||
self.assertEqual(q.pop().url, req1.url)
|
||||
self.assertEqual(len(q), 2)
|
||||
self.assertEqual(q.peek().url, req2.url)
|
||||
self.assertEqual(q.pop().url, req2.url)
|
||||
self.assertEqual(len(q), 1)
|
||||
self.assertEqual(q.peek().url, req3.url)
|
||||
self.assertEqual(q.pop().url, req3.url)
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.peek())
|
||||
self.assertIsNone(q.pop())
|
||||
q.close()
|
||||
|
||||
def test_fifo_without_peek(self):
|
||||
if hasattr(queuelib.queue.FifoMemoryQueue, "peek"):
|
||||
raise unittest.SkipTest("The queuelib queues do not define peek")
|
||||
q = self.queue()
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.pop())
|
||||
req1 = Request("http://www.example.com/1")
|
||||
req2 = Request("http://www.example.com/2")
|
||||
req3 = Request("http://www.example.com/3")
|
||||
q.push(req1)
|
||||
q.push(req2)
|
||||
q.push(req3)
|
||||
with self.assertRaises(NotImplementedError, msg="The underlying queue class does not implement 'peek'"):
|
||||
q.peek()
|
||||
self.assertEqual(len(q), 3)
|
||||
self.assertEqual(q.pop().url, req1.url)
|
||||
self.assertEqual(len(q), 2)
|
||||
self.assertEqual(q.pop().url, req2.url)
|
||||
self.assertEqual(len(q), 1)
|
||||
self.assertEqual(q.pop().url, req3.url)
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.pop())
|
||||
q.close()
|
||||
|
||||
|
||||
class LifoQueueMixin(RequestQueueTestMixin):
|
||||
def test_lifo_with_peek(self):
|
||||
if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"):
|
||||
raise unittest.SkipTest("The queuelib queues do not define peek")
|
||||
q = self.queue()
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.peek())
|
||||
self.assertIsNone(q.pop())
|
||||
req1 = Request("http://www.example.com/1")
|
||||
req2 = Request("http://www.example.com/2")
|
||||
req3 = Request("http://www.example.com/3")
|
||||
q.push(req1)
|
||||
q.push(req2)
|
||||
q.push(req3)
|
||||
self.assertEqual(len(q), 3)
|
||||
self.assertEqual(q.peek().url, req3.url)
|
||||
self.assertEqual(q.pop().url, req3.url)
|
||||
self.assertEqual(len(q), 2)
|
||||
self.assertEqual(q.peek().url, req2.url)
|
||||
self.assertEqual(q.pop().url, req2.url)
|
||||
self.assertEqual(len(q), 1)
|
||||
self.assertEqual(q.peek().url, req1.url)
|
||||
self.assertEqual(q.pop().url, req1.url)
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.peek())
|
||||
self.assertIsNone(q.pop())
|
||||
q.close()
|
||||
|
||||
def test_lifo_without_peek(self):
|
||||
if hasattr(queuelib.queue.FifoMemoryQueue, "peek"):
|
||||
raise unittest.SkipTest("The queuelib queues do not define peek")
|
||||
q = self.queue()
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.pop())
|
||||
req1 = Request("http://www.example.com/1")
|
||||
req2 = Request("http://www.example.com/2")
|
||||
req3 = Request("http://www.example.com/3")
|
||||
q.push(req1)
|
||||
q.push(req2)
|
||||
q.push(req3)
|
||||
with self.assertRaises(NotImplementedError, msg="The underlying queue class does not implement 'peek'"):
|
||||
q.peek()
|
||||
self.assertEqual(len(q), 3)
|
||||
self.assertEqual(q.pop().url, req3.url)
|
||||
self.assertEqual(len(q), 2)
|
||||
self.assertEqual(q.pop().url, req2.url)
|
||||
self.assertEqual(len(q), 1)
|
||||
self.assertEqual(q.pop().url, req1.url)
|
||||
self.assertEqual(len(q), 0)
|
||||
self.assertIsNone(q.pop())
|
||||
q.close()
|
||||
|
||||
|
||||
class PickleFifoDiskQueueRequestTest(FifoQueueMixin, BaseQueueTestCase):
|
||||
def queue(self):
|
||||
return PickleFifoDiskQueue.from_crawler(crawler=self.crawler, key="pickle/fifo")
|
||||
|
||||
|
||||
class PickleLifoDiskQueueRequestTest(LifoQueueMixin, BaseQueueTestCase):
|
||||
def queue(self):
|
||||
return PickleLifoDiskQueue.from_crawler(crawler=self.crawler, key="pickle/lifo")
|
||||
|
||||
|
||||
class MarshalFifoDiskQueueRequestTest(FifoQueueMixin, BaseQueueTestCase):
|
||||
def queue(self):
|
||||
return MarshalFifoDiskQueue.from_crawler(crawler=self.crawler, key="marshal/fifo")
|
||||
|
||||
|
||||
class MarshalLifoDiskQueueRequestTest(LifoQueueMixin, BaseQueueTestCase):
|
||||
def queue(self):
|
||||
return MarshalLifoDiskQueue.from_crawler(crawler=self.crawler, key="marshal/lifo")
|
||||
|
||||
|
||||
class FifoMemoryQueueRequestTest(FifoQueueMixin, BaseQueueTestCase):
|
||||
def queue(self):
|
||||
return FifoMemoryQueue.from_crawler(crawler=self.crawler)
|
||||
|
||||
|
||||
class LifoMemoryQueueRequestTest(LifoQueueMixin, BaseQueueTestCase):
|
||||
def queue(self):
|
||||
return LifoMemoryQueue.from_crawler(crawler=self.crawler)
|
||||
|
|
@ -108,7 +108,7 @@ class WarnWhenSubclassedTest(unittest.TestCase):
|
|||
|
||||
# ignore subclassing warnings
|
||||
with warnings.catch_warnings():
|
||||
warnings.simplefilter('ignore', ScrapyDeprecationWarning)
|
||||
warnings.simplefilter('ignore', MyWarning)
|
||||
|
||||
class UserClass(Deprecated):
|
||||
pass
|
||||
|
|
|
|||
|
|
@ -4,10 +4,11 @@ import operator
|
|||
import platform
|
||||
from datetime import datetime
|
||||
from itertools import count
|
||||
from warnings import catch_warnings
|
||||
from warnings import catch_warnings, filterwarnings
|
||||
|
||||
from twisted.trial import unittest
|
||||
|
||||
from scrapy.exceptions import ScrapyDeprecationWarning
|
||||
from scrapy.utils.asyncgen import as_async_generator, collect_asyncgen
|
||||
from scrapy.utils.defer import deferred_f_from_coro_f, aiter_errback
|
||||
from scrapy.utils.python import (
|
||||
|
|
@ -214,7 +215,11 @@ class UtilsPythonTestCase(unittest.TestCase):
|
|||
pass
|
||||
|
||||
_values = count()
|
||||
wk = WeakKeyCache(lambda k: next(_values))
|
||||
|
||||
with catch_warnings():
|
||||
filterwarnings("ignore", category=ScrapyDeprecationWarning)
|
||||
wk = WeakKeyCache(lambda k: next(_values))
|
||||
|
||||
k = _Weakme()
|
||||
v = wk[k]
|
||||
self.assertEqual(v, wk[k])
|
||||
|
|
|
|||
11
tox.ini
11
tox.ini
|
|
@ -50,6 +50,8 @@ commands =
|
|||
basepython = python3
|
||||
deps =
|
||||
{[testenv]deps}
|
||||
# Twisted[http2] is required to import some files
|
||||
Twisted[http2]>=17.9.0
|
||||
pytest-flake8
|
||||
commands =
|
||||
py.test --flake8 {posargs:docs scrapy tests}
|
||||
|
|
@ -57,12 +59,7 @@ commands =
|
|||
[testenv:pylint]
|
||||
basepython = python3
|
||||
deps =
|
||||
{[testenv]deps}
|
||||
# Optional dependencies
|
||||
boto
|
||||
reppy
|
||||
robotexclusionrulesparser
|
||||
# Test dependencies
|
||||
{[testenv:extra-deps]deps}
|
||||
pylint
|
||||
commands =
|
||||
pylint conftest.py docs extras scrapy setup.py tests
|
||||
|
|
@ -119,9 +116,11 @@ setenv =
|
|||
[testenv:extra-deps]
|
||||
deps =
|
||||
{[testenv]deps}
|
||||
boto
|
||||
reppy
|
||||
robotexclusionrulesparser
|
||||
Pillow>=4.0.0
|
||||
Twisted[http2]>=17.9.0
|
||||
|
||||
[testenv:asyncio]
|
||||
commands =
|
||||
|
|
|
|||
Loading…
Reference in New Issue