diff --git a/.github/workflows/tests-ubuntu.yml b/.github/workflows/tests-ubuntu.yml index 34819f227..06da46ca1 100644 --- a/.github/workflows/tests-ubuntu.yml +++ b/.github/workflows/tests-ubuntu.yml @@ -38,6 +38,9 @@ jobs: - python-version: pypy3.10 env: TOXENV: pypy3 + - python-version: pypy3.11 + env: + TOXENV: pypy3 # pinned deps - python-version: "3.9.21" @@ -59,7 +62,7 @@ jobs: - python-version: "3.13" env: TOXENV: extra-deps - - python-version: pypy3.10 + - python-version: pypy3.11 env: TOXENV: pypy3-extra-deps - python-version: "3.13" diff --git a/MANIFEST.in b/MANIFEST.in deleted file mode 100644 index 7700ae7bd..000000000 --- a/MANIFEST.in +++ /dev/null @@ -1,22 +0,0 @@ -include CODE_OF_CONDUCT.md -include CONTRIBUTING.md -include INSTALL.md -include NEWS -include SECURITY.md - -include scrapy/VERSION -include scrapy/mime.types -include scrapy/py.typed - -include codecov.yml -include conftest.py -include tox.ini - -recursive-include scrapy/templates * -recursive-include docs * -prune docs/build - -recursive-include extras * -recursive-include tests * - -global-exclude __pycache__ *.py[cod] diff --git a/docs/_templates/layout.html b/docs/_templates/layout.html new file mode 100644 index 000000000..6ec565e24 --- /dev/null +++ b/docs/_templates/layout.html @@ -0,0 +1,23 @@ +{% extends "!layout.html" %} + +{# Overriden to include a link to scrapy.org, not just to the docs root #} +{%- block sidebartitle %} + +{# the logo helper function was removed in Sphinx 6 and deprecated since Sphinx 4 #} +{# the master_doc variable was renamed to root_doc in Sphinx 4 (master_doc still exists in later Sphinx versions) #} +{%- set _logo_url = logo_url|default(pathto('_static/' + (logo or ""), 1)) %} +{%- set _root_doc = root_doc|default(master_doc) %} +scrapy.org / docs + +{%- if READTHEDOCS or DEBUG %} + {%- if theme_version_selector or theme_language_selector %} +
+
+
+
+ {%- endif %} +{%- endif %} + +{%- include "searchbox.html" %} + +{%- endblock %} diff --git a/docs/faq.rst b/docs/faq.rst index 79197d21d..95ae33a35 100644 --- a/docs/faq.rst +++ b/docs/faq.rst @@ -96,26 +96,14 @@ How can I simulate a user login in my spider? See :ref:`topics-request-response-ref-request-userlogin`. + .. _faq-bfo-dfo: Does Scrapy crawl in breadth-first or depth-first order? -------------------------------------------------------- -By default, Scrapy uses a `LIFO`_ queue for storing pending requests, which -basically means that it crawls in `DFO order`_. This order is more convenient -in most cases. +:ref:`DFO by default, but other orders are possible `. -To crawl in `BFO order`_: - -- Set :setting:`DEPTH_PRIORITY` to ``1``. - -- Set :setting:`SCHEDULER_MEMORY_QUEUE` to - :class:`~scrapy.squeues.FifoMemoryQueue`. - -- Set :setting:`SCHEDULER_DISK_QUEUE` to - :class:`~scrapy.squeues.PickleFifoDiskQueue`. - -- :ref:`Send start requests before other requests `. My Scrapy crawler has memory leaks. What can I do? -------------------------------------------------- @@ -431,6 +419,3 @@ See :issue:`2680`. .. _Python standard library modules: https://docs.python.org/3/py-modindex.html .. _Python package: https://pypi.org/ .. _user agents: https://en.wikipedia.org/wiki/User_agent -.. _LIFO: https://en.wikipedia.org/wiki/Stack_(abstract_data_type) -.. _DFO order: https://en.wikipedia.org/wiki/Depth-first_search -.. _BFO order: https://en.wikipedia.org/wiki/Breadth-first_search diff --git a/docs/news.rst b/docs/news.rst index edfa3b904..30d6b3b58 100644 --- a/docs/news.rst +++ b/docs/news.rst @@ -23,13 +23,19 @@ Modified requirements Backward-incompatible changes ~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ +- The ``from_settings()`` method of + :class:`~scrapy.spidermiddlewares.urllength.UrlLengthMiddleware` is removed + without a deprecation period (this was needed because after the + introduction of the + :class:`~scrapy.spidermiddlewares.base.BaseSpiderMiddleware` base class and + switching built-in spider middlewares to it those middlewares need the + :class:`~scrapy.crawler.Crawler` instance at run time). Please use + ``from_crawler()`` instead. + - The iteration of start requests and items no longer stops once there are requests in the scheduler, and instead runs continuously until all start requests have been scheduled. - As a result, the order in which start requests are sent may change. See - :ref:`start-requests` for details. - To restore the previous behavior, :ref:`use lazy start request scheduling `. @@ -63,6 +69,12 @@ Backward-incompatible changes instead of being defined as a generator, is now executed *after* the :ref:`scheduler ` instance has been created. +- When using :setting:`JOBDIR`, :ref:`start requests ` are + now serialized into their own, ``s``-suffixed priority folders. You can set + :setting:`SCHEDULER_START_DISK_QUEUE` to ``None`` or ``""`` to change that, + but the side effects may be undesirable. See + :setting:`SCHEDULER_START_DISK_QUEUE` for details. + Deprecations ~~~~~~~~~~~~ @@ -79,6 +91,12 @@ Deprecations (:issue:`456`, :issue:`3477`, :issue:`4467`, :issue:`5627`, :issue:`6729`) +- The ``__init__`` method of priority queue classes (see + :setting:`SCHEDULER_PRIORITY_QUEUE`) should now support a keyword-only + ``start_queue_cls`` parameter. + + (:issue:`6752`) + New features ~~~~~~~~~~~~ @@ -133,6 +151,13 @@ New features - Added new signals: :signal:`spider_start_blocking`, :signal:`scheduler_empty`. +- Added new settings: :setting:`SCHEDULER_START_DISK_QUEUE` and + :setting:`SCHEDULER_START_MEMORY_QUEUE`. + +- Added :class:`~scrapy.spidermiddlewares.start.StartSpiderMiddleware`, which + sets :reqmeta:`is_start_request` to ``True`` on :ref:`start requests + `. + - In :class:`~scrapy.core.engine.ExecutionEngine`: - Added a :attr:`~scrapy.core.engine.ExecutionEngine.scheduler` diff --git a/docs/topics/components.rst b/docs/topics/components.rst index 3a7644379..56f8c6498 100644 --- a/docs/topics/components.rst +++ b/docs/topics/components.rst @@ -37,6 +37,10 @@ That includes the classes that you may assign to the following settings: - :setting:`SCHEDULER_PRIORITY_QUEUE` +- :setting:`SCHEDULER_START_DISK_QUEUE` + +- :setting:`SCHEDULER_START_MEMORY_QUEUE` + - :setting:`SPIDER_MIDDLEWARES` Third-party Scrapy components may also let you define additional Scrapy diff --git a/docs/topics/coroutines.rst b/docs/topics/coroutines.rst index 18aa8c68a..448bf07e7 100644 --- a/docs/topics/coroutines.rst +++ b/docs/topics/coroutines.rst @@ -73,18 +73,103 @@ In addition to native coroutine APIs Scrapy has some APIs that return a :class:`~twisted.internet.defer.Deferred` object or take a user-supplied function that returns a :class:`~twisted.internet.defer.Deferred` object. These APIs are also asynchronous but don't yet support native ``async def`` syntax. -For example: +In the future we plan to add support for the ``async def`` syntax to these APIs +or replace them with other APIs where changing the existing ones is +possible. -- The :meth:`ExecutionEngine.download` method returns a - :class:`~twisted.internet.defer.Deferred` object. -- A custom download handler needs to define a ``download_request()`` method that - returns a :class:`~twisted.internet.defer.Deferred` object. +The following Scrapy methods return :class:`~twisted.internet.defer.Deferred` +objects (this list is not complete as it only includes methods that we think +may be useful for user code): + +- :class:`scrapy.crawler.Crawler`: + + - :meth:`~scrapy.crawler.Crawler.crawl` + + - :meth:`~scrapy.crawler.Crawler.stop` + +- :class:`scrapy.crawler.CrawlerRunner` (also inherited by + :class:`scrapy.crawler.CrawlerProcess`): + + - :meth:`~scrapy.crawler.CrawlerRunner.crawl` + + - :meth:`~scrapy.crawler.CrawlerRunner.stop` + + - :meth:`~scrapy.crawler.CrawlerRunner.join` + +- :class:`scrapy.core.engine.ExecutionEngine`: + + - :meth:`~scrapy.core.engine.ExecutionEngine.download` + +- :class:`scrapy.signalmanager.SignalManager`: + + - :meth:`~scrapy.signalmanager.SignalManager.send_catch_log_deferred` + +- :class:`~scrapy.mail.MailSender` + + - :meth:`~scrapy.mail.MailSender.send` + +The following user-supplied methods can return +:class:`~twisted.internet.defer.Deferred` objects (the methods that can also +return coroutines are listed in :ref:`coroutine-support`): + +- Custom download handlers (see :setting:`DOWNLOAD_HANDLERS`): + + - ``download_request()`` + + - ``close()`` + +- Custom downloader implementations (see :setting:`DOWNLOADER`): + + - ``fetch()`` + +- Custom scheduler implementations (see :setting:`SCHEDULER`): + + - :meth:`~scrapy.core.scheduler.BaseScheduler.open` + + - :meth:`~scrapy.core.scheduler.BaseScheduler.close` + +- Custom dupefilters (see :setting:`DUPEFILTER_CLASS`): + + - ``open()`` + + - ``close()`` + +- Custom feed storages (see :setting:`FEED_STORAGES`): + + - ``store()`` + +- Subclasses of :class:`scrapy.pipelines.media.MediaPipeline`: + + - ``media_to_download()`` + + - ``item_completed()`` + +- Custom storages used by subclasses of + :class:`scrapy.pipelines.files.FilesPipeline`: + + - ``persist_file()`` + + - ``stat_file()`` In most cases you can use these APIs in code that otherwise uses coroutines, by wrapping a :class:`~twisted.internet.defer.Deferred` object into a :class:`~asyncio.Future` object or vice versa. See :ref:`asyncio-await-dfd` for more information about this. +For example: + +- The :meth:`ExecutionEngine.download() + ` method returns a + :class:`~twisted.internet.defer.Deferred` object that fires with the + downloaded response. You can use this object directly in Deferred-based + code or convert it into a :class:`~asyncio.Future` object with + :func:`~scrapy.utils.defer.maybe_deferred_to_future`. +- A custom download handler needs to define a ``download_request()`` method + that returns a :class:`~twisted.internet.defer.Deferred` object. You can + write a method that works with Deferreds and returns one directly, or you + can write a coroutine and convert it into a function that returns a + Deferred with :func:`~scrapy.utils.defer.deferred_f_from_coro_f`. + General usage ============= diff --git a/docs/topics/feed-exports.rst b/docs/topics/feed-exports.rst index aff3cb29f..7b9aac8a8 100644 --- a/docs/topics/feed-exports.rst +++ b/docs/topics/feed-exports.rst @@ -517,18 +517,18 @@ as a fallback value if that key is not provided for a specific feed definition: FEED_EXPORT_ENCODING -------------------- -Default: ``None`` +Default: ``"utf-8"`` (:ref:`fallback `: ``None``) The encoding to be used for the feed. -If unset or set to ``None`` (default) it uses UTF-8 for everything except JSON output, -which uses safe numeric encoding (``\uXXXX`` sequences) for historic reasons. +If set to ``None``, it uses UTF-8 for everything except JSON output, which uses +safe numeric encoding (``\uXXXX`` sequences) for historic reasons. -Use ``utf-8`` if you want UTF-8 for JSON too. +Use ``"utf-8"`` if you want UTF-8 for JSON too. .. versionchanged:: 2.8 The :command:`startproject` command now sets this setting to - ``utf-8`` in the generated ``settings.py`` file. + ``"utf-8"`` in the generated ``settings.py`` file. .. setting:: FEED_EXPORT_FIELDS diff --git a/docs/topics/jobs.rst b/docs/topics/jobs.rst index 0e705dc64..50bcaa6d6 100644 --- a/docs/topics/jobs.rst +++ b/docs/topics/jobs.rst @@ -46,9 +46,9 @@ Keeping persistent state between batches Sometimes you'll want to keep some persistent spider state between pause/resume batches. You can use the ``spider.state`` attribute for that, which should be a -dict. There's :ref:`a built-in extension ` that takes care of serializing, storing and -loading that attribute from the job directory, when the spider starts and -stops. +dict. There's :ref:`a built-in extension ` +that takes care of serializing, storing and loading that attribute from the job +directory, when the spider starts and stops. Here's an example of a callback that uses the spider state (other spider code is omitted for brevity): diff --git a/docs/topics/request-response.rst b/docs/topics/request-response.rst index d00b11919..cd1a3b7f3 100644 --- a/docs/topics/request-response.rst +++ b/docs/topics/request-response.rst @@ -647,6 +647,7 @@ Those are: * ``ftp_user`` (See :setting:`FTP_USER` for more info) * :reqmeta:`handle_httpstatus_all` * :reqmeta:`handle_httpstatus_list` +* :reqmeta:`is_start_request` * :reqmeta:`max_retry_times` * :reqmeta:`proxy` * :reqmeta:`redirect_reasons` diff --git a/docs/topics/scheduler.rst b/docs/topics/scheduler.rst index 411ebe0f5..a24280cd7 100644 --- a/docs/topics/scheduler.rst +++ b/docs/topics/scheduler.rst @@ -54,8 +54,8 @@ Built-in scheduler ------------------ .. autoclass:: Scheduler() - :members: __len__, pause, unpause - :member-order: bysource + :members: + :special-members: __init__, __len__, pause, unpause .. _priority-queues: @@ -119,6 +119,7 @@ Writing a priority queue :special-members: __init__, __len__ :member-order: bysource + .. _custom-internal-queue: Writing an internal queue diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index 452d4b209..546861363 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -162,8 +162,17 @@ Those command-specific default settings are specified in the 6. Default global settings -------------------------- -The global defaults are located in the ``scrapy.settings.default_settings`` -module and documented in the :ref:`topics-settings-ref` section. +The ``scrapy.settings.default_settings`` module defines global default values +for some :ref:`built-in settings `. + +.. note:: :command:`startproject` generates a ``settings.py`` file that sets + some settings to different values. + + The reference documentation of settings indicates the default value if one + exists. If :command:`startproject` sets a value, that value is documented + as default, and the value from ``scrapy.settings.default_settings`` is + documented as “fallback”. + Compatibility with pickle ========================= @@ -461,7 +470,7 @@ Note that the event loop class must inherit from :class:`asyncio.AbstractEventLo BOT_NAME -------- -Default: ``'scrapybot'`` +Default: ```` (:ref:`fallback `: ``'scrapybot'``) The name of the bot implemented by this Scrapy project (also known as the project name). This name will be used for the logging too. @@ -1320,6 +1329,7 @@ Default: ``{}`` A dict containing the pipelines enabled by default in Scrapy. You should never modify this setting in your project, modify :setting:`ITEM_PIPELINES` instead. + .. setting:: JOBDIR JOBDIR @@ -1330,6 +1340,7 @@ Default: ``None`` A string indicating the directory for storing the state of a crawl when :ref:`pausing and resuming crawls `. + .. setting:: LOG_ENABLED LOG_ENABLED @@ -1566,7 +1577,7 @@ email notifying about it. If zero, no warning will be produced. NEWSPIDER_MODULE ---------------- -Default: ``''`` +Default: ``".spiders"`` (:ref:`fallback `: ``""``) Module where to create new spiders using the :command:`genspider` command. @@ -1625,9 +1636,7 @@ Adjust redirect request priority relative to original request: ROBOTSTXT_OBEY -------------- -Default: ``False`` - -Scope: ``scrapy.downloadermiddlewares.robotstxt`` +Default: ``True`` (:ref:`fallback `: ``False``) If enabled, Scrapy will respect robots.txt policies. For more information see :ref:`topics-dlmw-robots`. @@ -1737,6 +1746,51 @@ The default scheduler, :class:`~scrapy.core.scheduler.Scheduler`, uses this component. +.. setting:: SCHEDULER_START_DISK_QUEUE + +SCHEDULER_START_DISK_QUEUE +-------------------------- + +Default: :class:`~scrapy.squeues.PickleFifoDiskQueue` + +Type of disk queue (see :setting:`JOBDIR`) that the :ref:`scheduler +` uses for :ref:`start requests `. + +For available choices, see :ref:`disk-queues` and :ref:`custom-internal-queue`. + +.. queue-common-starts + +Use ``None`` or ``""`` to disable these separate queues entirely, and instead +have start requests share the same queues as other requests. + +.. note:: + + Disabling separate start request queues makes :ref:`start request order + ` unintuitive: start requests will be sent in order + only until :setting:`CONCURRENT_REQUESTS` is reached, then remaining start + requests will be sent in reverse order. + +.. queue-common-ends + + +.. setting:: SCHEDULER_START_MEMORY_QUEUE + +SCHEDULER_START_MEMORY_QUEUE +---------------------------- + +Default: :class:`~scrapy.squeues.FifoMemoryQueue` + +Type of in-memory queue that the :ref:`scheduler ` uses for +:ref:`start requests `. + +For available choices, see :ref:`memory-queues` and +:ref:`custom-internal-queue`. + +.. include:: settings.rst + :start-after: queue-common-starts + :end-before: queue-common-ends + + .. setting:: SCRAPER_SLOT_MAX_ACTIVE_SIZE SCRAPER_SLOT_MAX_ACTIVE_SIZE @@ -1856,7 +1910,7 @@ the spider. For more info see :ref:`topics-spider-middleware-setting`. SPIDER_MODULES -------------- -Default: ``[]`` +Default: ``[".spiders"]`` (:ref:`fallback `: ``[]``) A list of modules where Scrapy will look for spiders. diff --git a/docs/topics/spider-middleware.rst b/docs/topics/spider-middleware.rst index 760abd3f4..7f1855322 100644 --- a/docs/topics/spider-middleware.rst +++ b/docs/topics/spider-middleware.rst @@ -190,6 +190,19 @@ one or more of these methods: :param spider: the spider which raised the exception :type spider: :class:`~scrapy.Spider` object +Base class for custom spider middlewares +---------------------------------------- + +Scrapy provides a base class for custom spider middlewares. It's not required +to use it but it can help with simplifying middleware implementations and +reducing the amount of boilerplate code in :ref:`universal middlewares +`. + +.. module:: scrapy.spidermiddlewares.base + +.. autoclass:: BaseSpiderMiddleware + :members: + .. _topics-spider-middleware-ref: Built-in spider middleware reference @@ -405,6 +418,14 @@ String value Class name (as a string) .. _"unsafe-url": https://www.w3.org/TR/referrer-policy/#referrer-policy-unsafe-url +StartSpiderMiddleware +--------------------- + +.. module:: scrapy.spidermiddlewares.start + +.. autoclass:: StartSpiderMiddleware + + UrlLengthMiddleware ------------------- diff --git a/docs/topics/spiders.rst b/docs/topics/spiders.rst index 107fd935a..5d527c59c 100644 --- a/docs/topics/spiders.rst +++ b/docs/topics/spiders.rst @@ -369,81 +369,12 @@ See `Scrapyd documentation`_. Start requests ============== -**Start requests** are the :ref:`requests ` that serve as a starting -point for a spider. +**Start requests** are :class:`~scrapy.Request` objects yielded from the +:meth:`~scrapy.Spider.start` method of a spider or from the +:meth:`~scrapy.spidermiddlewares.SpiderMiddleware.process_start` method of a +:ref:`spider middleware `. -They are yielded by the :meth:`~scrapy.Spider.start` method, and may be -modified by :ref:`spider middlewares `. - -They are not necessarily the *first* requests sent, and may not be sent in -order; reaching :setting:`CONCURRENT_REQUESTS` and :ref:`scheduling -` start requests is prioritized by default. But you can -change that. - -.. - The request send order when all start requests and callback requests have - the same priority is rather unintuitive: - - 1. First, the first CONCURRENT_REQUESTS start requests are sent in order. - - Awaiting slow operations in Spider.start() can lower that. - - 2. Then, assuming an even domain distribution in start requests (i.e. - ABCABC, not AABBCC), the last N start requests are sent in reverse - order, where N is: - - min(CONCURRENT_REQUESTS, CONCURRENT_REQUESTS_PER_DOMAIN * domain_count) - - 3. Finally, the remaining start requests are also sent in reverse order, - but only when there are not enough pending requests yielded from - callbacks to reach the configured concurrency. - - The reverse order is because the scheduler uses a LIFO queue by default - (SCHEDULER_MEMORY_QUEUE, SCHEDULER_DISK_QUEUE). The order of the first few - requests is unnaffected because they are sent as soon as they are - scheduled. The last start requests sent before callback requests are those - that can be sent before the first callback requests are scheduled. - - We do not document this behavior, so that we may change it in the future - without breaking the contract. - -.. _start-requests-order: - -Start request order -------------------- - -To force a specific request order, override the :meth:`~scrapy.Spider.start` -method to set :attr:`Request.priority `. - -For example: - -- To send start requests before other requests: - - .. code-block:: python - - async def start(self): - async for item_or_request in super().start(): - if isinstance(item_or_request, Request): - item_or_request = item_or_request.replace(priority=1) - yield item_or_request - -- To send start requests in yield order: - - .. code-block:: python - - async def start(self): - priority = len(self.start_urls) - async for item_or_request in super().start(): - if isinstance(item_or_request, Request): - item_or_request = item_or_request.replace(priority=priority) - yield item_or_request - priority -= 1 - -You can also :ref:`customize the scheduler ` if you need -more control over request prioritization. - -.. note:: By default, :ref:`the first few requests are sent in yield order - ` regardless, for performance reasons. +.. seealso:: :ref:`start-request-order` .. _start-requests-front-load: @@ -452,7 +383,7 @@ Start request front loading By default, the first few requests yielded by :meth:`~scrapy.Spider.start` are sent in yield order, even if you :ref:`try to sort them differently -`. This is because, until :setting:`CONCURRENT_REQUESTS` +`. This is because, until :setting:`CONCURRENT_REQUESTS` is reached, start requests are removed from the scheduler immediately after they are scheduled. diff --git a/pyproject.toml b/pyproject.toml index ddafa2a41..c66e940d2 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [build-system] -requires = ["setuptools >= 61.0"] -build-backend = "setuptools.build_meta" +requires = ["hatchling>=1.27.0"] +build-backend = "hatchling.build" [project] name = "Scrapy" @@ -10,29 +10,28 @@ dependencies = [ "Twisted>=21.7.0", "cryptography>=37.0.0", "cssselect>=0.9.1", + "defusedxml>=0.7.1", + "itemadapter>=0.1.0", "itemloaders>=1.0.1", + "lxml>=4.9.3", + "packaging", "parsel>=1.5.0", + "protego>=0.1.15", "pyOpenSSL>=22.0.0", "queuelib>=1.4.2", "service_identity>=18.1.0", + "tldextract", "w3lib>=1.17.0", "zope.interface>=5.1.0", - "protego>=0.1.15", - "itemadapter>=0.1.0", - "packaging", - "tldextract", - "lxml>=4.9.3", - "defusedxml>=0.7.1", # Platform-specific dependencies 'PyDispatcher>=2.0.5; platform_python_implementation == "CPython"', 'PyPyDispatcher>=2.1.0; platform_python_implementation == "PyPy"', ] classifiers = [ - "Framework :: Scrapy", "Development Status :: 5 - Production/Stable", "Environment :: Console", + "Framework :: Scrapy", "Intended Audience :: Developers", - "License :: OSI Approved :: BSD License", "Operating System :: OS Independent", "Programming Language :: Python", "Programming Language :: Python :: 3", @@ -47,6 +46,8 @@ classifiers = [ "Topic :: Software Development :: Libraries :: Application Frameworks", "Topic :: Software Development :: Libraries :: Python Modules", ] +license = "BSD-3-Clause" +license-files = ["LICENSE", "AUTHORS"] readme = "README.rst" requires-python = ">=3.9" authors = [{ name = "Scrapy developers", email = "pablo@pablohoffman.com" }] @@ -63,12 +64,26 @@ releasenotes = "https://docs.scrapy.org/en/latest/news.html" [project.scripts] scrapy = "scrapy.cmdline:execute" -[tool.setuptools.packages.find] -where = ["."] -include = ["scrapy", "scrapy.*",] +[tool.hatch.build.targets.sdist] +include = [ + "/docs", + "/extras", + "/scrapy", + "/tests", + "/tests_typing", + "/CODE_OF_CONDUCT.md", + "/CONTRIBUTING.md", + "/INSTALL.md", + "/NEWS", + "/SECURITY.md", + "/codecov.yml", + "/conftest.py", + "/tox.ini", +] -[tool.setuptools.dynamic] -version = {file = "./scrapy/VERSION"} +[tool.hatch.version] +path = "scrapy/VERSION" +pattern = "^(?P.+)$" [tool.mypy] ignore_missing_imports = true diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index b01b49d68..ccb227d24 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -12,12 +12,12 @@ from time import time from traceback import format_exc from typing import TYPE_CHECKING, Any, TypeVar, cast -from twisted.internet.defer import Deferred, succeed +from twisted.internet.defer import Deferred, inlineCallbacks, succeed from twisted.python.failure import Failure from scrapy import signals from scrapy.core.scheduler import BaseScheduler -from scrapy.core.scraper import Scraper, _HandleOutputDeferred +from scrapy.core.scraper import Scraper from scrapy.exceptions import CloseSpider, DontCloseSpider, IgnoreRequest from scrapy.http import Request, Response from scrapy.utils.defer import ( @@ -30,7 +30,7 @@ from scrapy.utils.python import global_object_name from scrapy.utils.reactor import CallLaterOnce if TYPE_CHECKING: - from collections.abc import AsyncIterator, Callable + from collections.abc import AsyncIterator, Callable, Generator from scrapy.core.downloader import Downloader from scrapy.crawler import Crawler @@ -344,11 +344,10 @@ class ExecutionEngine: ) return True + @inlineCallbacks def _handle_downloader_output( self, result: Request | Response | Failure, request: Request - ) -> _HandleOutputDeferred | None: - assert self.spider is not None # typing - + ) -> Generator[Deferred[Any], Any, None]: if not isinstance(result, (Request, Response, Failure)): raise TypeError( f"Incorrect type: expected Request, Response or Failure, got {type(result)}: {result!r}" @@ -357,17 +356,17 @@ class ExecutionEngine: # downloader middleware can return requests (for example, redirects) if isinstance(result, Request): self.crawl(result) - return None + return - 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), + try: + yield self.scraper.enqueue_scrape(result, request) + except Exception: + assert self.spider is not None + logger.error( + "Error while enqueuing scrape", + exc_info=True, extra={"spider": self.spider}, ) - ) - return d def spider_is_idle(self) -> bool: if self._slot is None: @@ -387,20 +386,20 @@ class ExecutionEngine: """Inject the request into the spider <-> downloader pipeline""" if self.spider is None: raise RuntimeError(f"No open spider to crawl: {request}") - self._schedule_request(request, self.spider) + self._schedule_request(request) self._slot.nextcall.schedule() # type: ignore[union-attr] - def _schedule_request(self, request: Request, spider: Spider) -> None: - assert self.scheduler is not None # typing + def _schedule_request(self, request: Request) -> None: request_scheduled_result = self.signals.send_catch_log( signals.request_scheduled, request=request, - spider=spider, + spider=self.spider, dont_log=IgnoreRequest, ) for handler, result in request_scheduled_result: if isinstance(result, Failure) and isinstance(result.value, IgnoreRequest): return + assert self.scheduler is not None try: request_was_enqueued = self.scheduler.enqueue_request(request) except Exception as exception: @@ -411,7 +410,7 @@ class ExecutionEngine: request_was_enqueued = False if not request_was_enqueued: self.signals.send_catch_log( - signals.request_dropped, request=request, spider=spider + signals.request_dropped, request=request, spider=self.spider ) def download(self, request: Request) -> Deferred[Response]: @@ -481,9 +480,7 @@ class ExecutionEngine: nextcall = CallLaterOnce(self._start_scheduled_requests) self.scheduler = build_from_crawler(self.scheduler_cls, self.crawler) self._slot = _Slot(close_if_idle, nextcall) - self._start = await maybe_deferred_to_future( - self.scraper.spidermw.process_start(spider) - ) + self._start = await self.scraper.spidermw.process_start(spider) if hasattr(self.scheduler, "open") and (d := self.scheduler.open(spider)): await maybe_deferred_to_future(d) await maybe_deferred_to_future(self.scraper.open_spider(spider)) @@ -543,7 +540,7 @@ class ExecutionEngine: dfd.addBoth(lambda _: self.downloader.close()) dfd.addErrback(log_failure("Downloader close failure")) - dfd.addBoth(lambda _: self.scraper.close_spider(spider)) + dfd.addBoth(lambda _: self.scraper.close_spider()) dfd.addErrback(log_failure("Scraper close failure")) if hasattr(self.scheduler, "close"): diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index 2dec0d732..3dfd709bf 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -5,13 +5,16 @@ import logging from abc import abstractmethod from pathlib import Path from typing import TYPE_CHECKING, Any, cast +from warnings import warn # working around https://github.com/sphinx-doc/sphinx/issues/10400 from twisted.internet.defer import Deferred # noqa: TC002 +from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.spiders import Spider # noqa: TC001 from scrapy.utils.job import job_dir from scrapy.utils.misc import build_from_crawler, load_object +from scrapy.utils.python import global_object_name if TYPE_CHECKING: # requires queuelib >= 1.6.2 @@ -117,13 +120,101 @@ class BaseScheduler(metaclass=BaseSchedulerMeta): class Scheduler(BaseScheduler): """Default :ref:`scheduler `. - Requests are stored in memory by default. Set :setting:`JOBDIR` to switch - to disk storage. + Requests are stored into priority queues + (:setting:`SCHEDULER_PRIORITY_QUEUE`) that sort requests by + :attr:`~scrapy.http.Request.priority`. - Requests are dropped if :attr:`~scrapy.Request.dont_filter` is ``False`` - and :setting:`DUPEFILTER_CLASS` flags them as duplicate requests. + By default, a single, memory-based priority queue is used for all requests. + When using :setting:`JOBDIR`, a disk-based priority queue is also created, + and only unserializable requests are stored in the memory-based priority + queue. For a given priority value, requests in memory take precedence over + requests in disk. - :setting:`SCHEDULER_PRIORITY_QUEUE` handles request prioritization. + Each priority queue stores requests in separate internal queues, one per + priority value. The memory priority queue uses + :setting:`SCHEDULER_MEMORY_QUEUE` queues, while the disk priority queue + uses :setting:`SCHEDULER_DISK_QUEUE` queues. The internal queues determine + :ref:`request order ` when requests have the same priority. + :ref:`Start requests ` are stored into separate internal + queues by default, and :ref:`ordered differently `. + + Duplicate requests are filtered out with an instance of + :setting:`DUPEFILTER_CLASS` if :attr:`~scrapy.Request.dont_filter` is + ``False``. + + .. seealso:: :ref:`topics-jobs` + + .. _request-order: + + Request order + ============= + + With default settings, pending requests are stored in a LIFO_ queue + (:ref:`except for start requests `). As a result, + crawling happens in `DFO order`_, which is usually the most convenient + crawl order. However, you can enforce :ref:`BFO ` or :ref:`a custom + order ` (:ref:`except for the first few requests + `). + + .. _LIFO: https://en.wikipedia.org/wiki/Stack_(abstract_data_type) + .. _DFO order: https://en.wikipedia.org/wiki/Depth-first_search + + .. _start-request-order: + + Start request order + ------------------- + + :ref:`Start requests ` are sent in the order they are + yielded from :meth:`~scrapy.Spider.start`, and given the same + :attr:`~scrapy.http.Request.priority`, start requests take precedence over + other requests. + + You can set :setting:`SCHEDULER_START_MEMORY_QUEUE` and + :setting:`SCHEDULER_START_DISK_QUEUE` to ``None`` to handle start requests + the same as other requests when it comes to order and priority. + + + .. _bfo: + + Crawling in BFO order + --------------------- + + If you do want to crawl in `BFO order`_, you can do it by setting the + following :ref:`settings `: + + | :setting:`DEPTH_PRIORITY` = ``1`` + | :setting:`SCHEDULER_DISK_QUEUE` = ``"scrapy.squeues.PickleFifoDiskQueue"`` + | :setting:`SCHEDULER_MEMORY_QUEUE` = ``"scrapy.squeues.FifoMemoryQueue"`` + + .. _BFO order: https://en.wikipedia.org/wiki/Breadth-first_search + + + .. _custom-request-order: + + Crawling in a custom order + -------------------------- + + You can manually set :attr:`~scrapy.http.Request.priority` on requests to + force a specific request order. + + + .. _concurrency-v-order: + + Concurrency affects order + ------------------------- + + While pending requests are below the configured values of + :setting:`CONCURRENT_REQUESTS`, :setting:`CONCURRENT_REQUESTS_PER_DOMAIN` + or :setting:`CONCURRENT_REQUESTS_PER_IP`, those requests are sent + concurrently. + + As a result, the first few requests of a crawl may not follow the desired + order. Lowering those settings to ``1`` enforces the desired order except + for the very first request, but it significantly slows down the crawl as a + whole. + + Stats + ===== The following stats are generated: @@ -141,31 +232,8 @@ class Scheduler(BaseScheduler): enabling :setting:`SCHEDULER_DEBUG` to log a warning message with details about the first unserializable request, to try and figure out how to make it serializable. - - .. seealso:: :ref:`topics-jobs` """ - def __init__( - self, - dupefilter: BaseDupeFilter, - jobdir: str | None = None, - dqclass: type[BaseQueue] | None = None, - mqclass: type[BaseQueue] | None = None, - logunser: bool = False, - stats: StatsCollector | None = None, - pqclass: type[PriorityQueueProtocol] | None = None, - crawler: Crawler | None = None, - ): - self.df: BaseDupeFilter = dupefilter - self.dqdir: str | None = self._dqdir(jobdir) - self.pqclass: type[PriorityQueueProtocol] | None = pqclass - self.dqclass: type[BaseQueue] | None = dqclass - self.mqclass: type[BaseQueue] | None = mqclass - self.logunser: bool = logunser - self.stats: StatsCollector | None = stats - self.crawler: Crawler | None = crawler - self._paused = False - @classmethod def from_crawler(cls, crawler: Crawler) -> Self: dupefilter_cls = load_object(crawler.settings["DUPEFILTER_CLASS"]) @@ -180,6 +248,79 @@ class Scheduler(BaseScheduler): crawler=crawler, ) + def __init__( + self, + dupefilter: BaseDupeFilter, + jobdir: str | None = None, + dqclass: type[BaseQueue] | None = None, + mqclass: type[BaseQueue] | None = None, + logunser: bool = False, + stats: StatsCollector | None = None, + pqclass: type[PriorityQueueProtocol] | None = None, + crawler: Crawler | None = None, + ): + """Initialize the scheduler. + + :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` + """ + self.df: BaseDupeFilter = dupefilter + self.dqdir: str | None = self._dqdir(jobdir) + self.pqclass: type[PriorityQueueProtocol] | None = pqclass + self.dqclass: type[BaseQueue] | None = dqclass + self.mqclass: type[BaseQueue] | None = mqclass + self.logunser: bool = logunser + self.stats: StatsCollector | None = stats + self.crawler: Crawler | None = crawler + self._paused = False + self._sdqclass: type[BaseQueue] | None = self._get_start_queue_cls( + crawler, "DISK" + ) + self._smqclass: type[BaseQueue] | None = self._get_start_queue_cls( + crawler, "MEMORY" + ) + + def _get_start_queue_cls( + self, crawler: Crawler | None, queue: str + ) -> type[BaseQueue] | None: + if crawler is None: + return None + cls = crawler.settings[f"SCHEDULER_START_{queue}_QUEUE"] + if not cls: + return None + return load_object(cls) + def has_pending_requests(self) -> bool: return len(self) > 0 @@ -278,12 +419,27 @@ class Scheduler(BaseScheduler): """Create a new priority queue instance, with in-memory storage""" assert self.crawler assert self.pqclass - return build_from_crawler( - self.pqclass, - self.crawler, - downstream_queue_cls=self.mqclass, - key="", - ) + try: + return build_from_crawler( + self.pqclass, + self.crawler, + downstream_queue_cls=self.mqclass, + key="", + start_queue_cls=self._smqclass, + ) + except TypeError: + warn( + f"The __init__ method of {global_object_name(self.pqclass)} " + f"does not support a `start_queue_cls` keyword-only " + f"parameter.", + ScrapyDeprecationWarning, + ) + return build_from_crawler( + self.pqclass, + self.crawler, + downstream_queue_cls=self.mqclass, + key="", + ) def _dq(self) -> PriorityQueueProtocol: """Create a new priority queue instance, with disk storage""" @@ -291,13 +447,29 @@ class Scheduler(BaseScheduler): assert self.dqdir assert self.pqclass state = self._read_dqs_state(self.dqdir) - q = build_from_crawler( - self.pqclass, - self.crawler, - downstream_queue_cls=self.dqclass, - key=self.dqdir, - startprios=state, - ) + try: + q = build_from_crawler( + self.pqclass, + self.crawler, + downstream_queue_cls=self.dqclass, + key=self.dqdir, + startprios=state, + start_queue_cls=self._sdqclass, + ) + except TypeError: + warn( + f"The __init__ method of {global_object_name(self.pqclass)} " + f"does not support a `start_queue_cls` keyword-only " + f"parameter.", + ScrapyDeprecationWarning, + ) + q = build_from_crawler( + self.pqclass, + self.crawler, + downstream_queue_cls=self.dqclass, + key=self.dqdir, + startprios=state, + ) if q: logger.info( "Resuming crawl (%(queuesize)d requests scheduled)", diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index 85dd9b559..9378f2651 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -4,22 +4,30 @@ extracts information from them""" from __future__ import annotations import logging +import warnings from collections import deque -from collections.abc import AsyncIterator, Iterator -from typing import TYPE_CHECKING, Any, TypeVar, Union, cast +from collections.abc import AsyncIterator +from typing import TYPE_CHECKING, Any, TypeVar, Union -from twisted.internet.defer import Deferred, inlineCallbacks +from twisted.internet.defer import Deferred, inlineCallbacks, maybeDeferred from twisted.python.failure import Failure from scrapy import Spider, signals from scrapy.core.spidermw import SpiderMiddlewareManager -from scrapy.exceptions import CloseSpider, DropItem, IgnoreRequest +from scrapy.exceptions import ( + CloseSpider, + DropItem, + IgnoreRequest, + ScrapyDeprecationWarning, +) from scrapy.http import Request, Response from scrapy.utils.defer import ( + _defer_sleep, aiter_errback, - defer_fail, - defer_succeed, + deferred_f_from_coro_f, + deferred_from_coro, iter_errback, + maybe_deferred_to_future, parallel, parallel_async, ) @@ -40,9 +48,7 @@ logger = logging.getLogger(__name__) _T = TypeVar("_T") -_ParallelResult = list[tuple[bool, Iterator[Any]]] -_HandleOutputDeferred = Deferred[Union[_ParallelResult, None]] -QueueTuple = tuple[Union[Response, Failure], Request, _HandleOutputDeferred] +QueueTuple = tuple[Union[Response, Failure], Request, Deferred[None]] class Slot: @@ -60,8 +66,9 @@ class Slot: def add_response_request( self, result: Response | Failure, request: Request - ) -> _HandleOutputDeferred: - deferred: _HandleOutputDeferred = Deferred() + ) -> Deferred[None]: + # this Deferred will be awaited in enqueue_scrape() + deferred: Deferred[None] = Deferred() self.queue.append((result, request, deferred)) if isinstance(result, Response): self.active_size += max(len(result.body), self.MIN_RESPONSE_SIZE) @@ -70,9 +77,9 @@ class Slot: return deferred def next_response_request_deferred(self) -> QueueTuple: - response, request, deferred = self.queue.popleft() + result, request, deferred = self.queue.popleft() self.active.add(request) - return response, request, deferred + return result, request, deferred def finish_response(self, result: Response | Failure, request: Request) -> None: self.active.remove(request) @@ -110,145 +117,192 @@ class Scraper: self.slot = Slot(self.crawler.settings.getint("SCRAPER_SLOT_MAX_ACTIVE_SIZE")) yield self.itemproc.open_spider(spider) - def close_spider(self, spider: Spider) -> Deferred[Spider]: + def close_spider(self, spider: Spider | None = None) -> Deferred[Spider]: """Close a spider being scraped and release its resources""" + if spider is not None: + warnings.warn( + "Passing a 'spider' argument to Scraper.close_spider() is deprecated.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + if self.slot is None: raise RuntimeError("Scraper slot not assigned") self.slot.closing = Deferred() self.slot.closing.addCallback(self.itemproc.close_spider) - self._check_if_closing(spider) + self._check_if_closing() return self.slot.closing def is_idle(self) -> bool: """Return True if there isn't any more spiders to process""" return not self.slot - def _check_if_closing(self, spider: Spider) -> None: + def _check_if_closing(self) -> None: assert self.slot is not None # typing + assert self.crawler.spider if self.slot.closing and self.slot.is_idle(): - self.slot.closing.callback(spider) + assert self.crawler.spider + self.slot.closing.callback(self.crawler.spider) + @inlineCallbacks def enqueue_scrape( - self, result: Response | Failure, request: Request, spider: Spider - ) -> _HandleOutputDeferred: + self, result: Response | Failure, request: Request, spider: Spider | None = None + ) -> Generator[Deferred[Any], Any, None]: + if spider is not None: + warnings.warn( + "Passing a 'spider' argument to Scraper.enqueue_scrape() is deprecated.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + if self.slot is None: raise RuntimeError("Scraper slot not assigned") dfd = self.slot.add_response_request(result, request) - - def finish_scraping(_: _T) -> _T: - assert self.slot is not None - self.slot.finish_response(result, request) - self._check_if_closing(spider) - self._scrape_next(spider) - return _ - - dfd.addBoth(finish_scraping) - dfd.addErrback( - lambda f: logger.error( + self._scrape_next() + try: + yield dfd + except Exception: + logger.error( "Scraper bug processing %(request)s", {"request": request}, - exc_info=failure_to_exc_info(f), - extra={"spider": spider}, + exc_info=True, + extra={"spider": self.crawler.spider}, ) - ) - self._scrape_next(spider) - return dfd + finally: + self.slot.finish_response(result, request) + self._check_if_closing() + self._scrape_next() - def _scrape_next(self, spider: Spider) -> None: + def _scrape_next(self) -> None: assert self.slot is not None # typing while self.slot.queue: - response, request, deferred = self.slot.next_response_request_deferred() - self._scrape(response, request, spider).chainDeferred(deferred) + result, request, deferred = self.slot.next_response_request_deferred() + self._scrape(result, request).chainDeferred(deferred) - def _scrape( - self, result: Response | Failure, request: Request, spider: Spider - ) -> _HandleOutputDeferred: - """ - Handle the downloaded response or failure through the spider callback/errback - """ + @deferred_f_from_coro_f + async def _scrape(self, result: Response | Failure, request: Request) -> None: + """Handle the downloaded response or failure through the spider callback/errback.""" if not isinstance(result, (Response, Failure)): raise TypeError( f"Incorrect type: expected Response or Failure, got {type(result)}: {result!r}" ) - dfd: Deferred[Iterable[Any] | AsyncIterator[Any]] = self._scrape2( - result, request, spider - ) # returns spider's processed output - dfd.addErrback(self.handle_spider_error, request, result, spider) - dfd2: _HandleOutputDeferred = dfd.addCallback( - self.handle_spider_output, request, cast(Response, result), spider - ) - return dfd2 - def _scrape2( - self, result: Response | Failure, request: Request, spider: Spider - ) -> Deferred[Iterable[Any] | AsyncIterator[Any]]: - """ - Handle the different cases of request's result been a Response or a Failure - """ + assert self.crawler.spider if isinstance(result, Response): - # Deferreds are invariant so Mutable*Chain isn't matched to *Iterable - return self.spidermw.scrape_response( # type: ignore[return-value] - self.call_spider, result, request, spider - ) - # else result is a Failure - dfd = self.call_spider(result, request, spider) - dfd.addErrback(self._log_download_errors, result, request, spider) - return dfd + try: + # call the spider middlewares and the request callback with the response + output = await maybe_deferred_to_future( + self.spidermw.scrape_response( + self.call_spider, result, request, self.crawler.spider + ) + ) + except Exception: + self.handle_spider_error(Failure(), request, result) + else: + await self.handle_spider_output_async(output, request, result) + return + + try: + # call the request errback with the downloader error + await self.call_spider_async(result, request) + except Exception as spider_exc: + # the errback didn't silence the exception + if not result.check(IgnoreRequest): + logkws = self.logformatter.download_error( + result, request, self.crawler.spider + ) + logger.log( + *logformatter_adapter(logkws), + extra={"spider": self.crawler.spider}, + exc_info=failure_to_exc_info(result), + ) + if spider_exc is not result.value: + # the errback raised a different exception, handle it + self.handle_spider_error(Failure(), request, result) def call_spider( - self, result: Response | Failure, request: Request, spider: Spider + self, result: Response | Failure, request: Request, spider: Spider | None = None ) -> Deferred[Iterable[Any] | AsyncIterator[Any]]: - dfd: Deferred[Any] + if spider is not None: + warnings.warn( + "Passing a 'spider' argument to Scraper.call_spider() is deprecated.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + return deferred_from_coro(self.call_spider_async(result, request)) + + async def call_spider_async( + self, result: Response | Failure, request: Request + ) -> Iterable[Any] | AsyncIterator[Any]: + """Call the request callback or errback with the response or failure.""" + await maybe_deferred_to_future(_defer_sleep()) + assert self.crawler.spider if isinstance(result, Response): if getattr(result, "request", None) is None: result.request = request assert result.request - callback = result.request.callback or spider._parse - warn_on_generator_with_return_value(spider, callback) - dfd = defer_succeed(result) - dfd.addCallbacks( - callback=callback, callbackKeywords=result.request.cb_kwargs - ) + callback = result.request.callback or self.crawler.spider._parse + warn_on_generator_with_return_value(self.crawler.spider, callback) + output = callback(result, **result.request.cb_kwargs) else: # result is a Failure # TODO: properly type adding this attribute to a Failure result.request = request # type: ignore[attr-defined] - dfd = defer_fail(result) - if request.errback: - warn_on_generator_with_return_value(spider, request.errback) - dfd.addErrback(request.errback) - dfd2: Deferred[Iterable[Any] | AsyncIterator[Any]] = dfd.addCallback( - iterate_spider_output + if not request.errback: + result.raiseException() + warn_on_generator_with_return_value(self.crawler.spider, request.errback) + output = request.errback(result) + if isinstance(output, Failure): + output.raiseException() + # else the errback returned actual output (like a callback), + # which needs to be passed to iterate_spider_output() + return await maybe_deferred_to_future( + maybeDeferred(iterate_spider_output, output) ) - return dfd2 def handle_spider_error( self, _failure: Failure, request: Request, response: Response | Failure, - spider: Spider, + spider: Spider | None = None, ) -> None: + """Handle an exception raised by a spider callback or errback.""" + if spider is not None: + warnings.warn( + "Passing a 'spider' argument to Scraper.handle_spider_error() is deprecated.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + + assert self.crawler.spider exc = _failure.value if isinstance(exc, CloseSpider): assert self.crawler.engine is not None # typing - self.crawler.engine.close_spider(spider, exc.reason or "cancelled") + self.crawler.engine.close_spider( + self.crawler.spider, exc.reason or "cancelled" + ) return - logkws = self.logformatter.spider_error(_failure, request, response, spider) + logkws = self.logformatter.spider_error( + _failure, request, response, self.crawler.spider + ) logger.log( *logformatter_adapter(logkws), exc_info=failure_to_exc_info(_failure), - extra={"spider": spider}, + extra={"spider": self.crawler.spider}, ) self.signals.send_catch_log( signal=signals.spider_error, failure=_failure, response=response, - spider=spider, + spider=self.crawler.spider, ) assert self.crawler.stats - self.crawler.stats.inc_value("spider_exceptions/count", spider=spider) self.crawler.stats.inc_value( - f"spider_exceptions/{_failure.value.__class__.__name__}", spider=spider + "spider_exceptions/count", spider=self.crawler.spider + ) + self.crawler.stats.inc_value( + f"spider_exceptions/{_failure.value.__class__.__name__}", + spider=self.crawler.spider, ) def handle_spider_output( @@ -256,55 +310,65 @@ class Scraper: result: Iterable[_T] | AsyncIterator[_T], request: Request, response: Response, - spider: Spider, - ) -> _HandleOutputDeferred: - if not result: - return defer_succeed(None) - it: Iterable[_T] | AsyncIterator[_T] - dfd: Deferred[_ParallelResult] - if isinstance(result, AsyncIterator): - it = aiter_errback( - result, self.handle_spider_error, request, response, spider + spider: Spider | None = None, + ) -> Deferred[None]: + """Pass items/requests produced by a callback to ``_process_spidermw_output()`` in parallel.""" + if spider is not None: + warnings.warn( + "Passing a 'spider' argument to Scraper.handle_spider_output() is deprecated.", + category=ScrapyDeprecationWarning, + stacklevel=2, ) - dfd = parallel_async( - it, - self.concurrent_items, - self._process_spidermw_output, - request, - response, - spider, - ) - else: - it = iter_errback( - result, self.handle_spider_error, request, response, spider - ) - dfd = parallel( - it, - self.concurrent_items, - self._process_spidermw_output, - request, - response, - spider, - ) - # returning Deferred[_ParallelResult] instead of Deferred[Union[_ParallelResult, None]] - return dfd # type: ignore[return-value] + return deferred_from_coro( + self.handle_spider_output_async(result, request, response) + ) - def _process_spidermw_output( - self, output: Any, request: Request, response: Response, spider: Spider - ) -> Deferred[Any] | None: + async def handle_spider_output_async( + self, + result: Iterable[_T] | AsyncIterator[_T], + request: Request, + response: Response, + ) -> None: + """Pass items/requests produced by a callback to ``_process_spidermw_output()`` in parallel.""" + if isinstance(result, AsyncIterator): + ait = aiter_errback(result, self.handle_spider_error, request, response) + await maybe_deferred_to_future( + parallel_async( + ait, + self.concurrent_items, + self._process_spidermw_output, + response, + ) + ) + return + it = iter_errback(result, self.handle_spider_error, request, response) + await maybe_deferred_to_future( + parallel( + it, + self.concurrent_items, + self._process_spidermw_output, + response, + ) + ) + + @deferred_f_from_coro_f + async def _process_spidermw_output(self, output: Any, response: Response) -> None: """Process each Request/Item (given in the output parameter) returned - from the given spider + from the given spider. + + Items are sent to the item pipelines, requests are scheduled. """ if isinstance(output, Request): assert self.crawler.engine is not None # typing self.crawler.engine.crawl(request=output) - elif output is None: - pass - else: - return self.start_itemproc(output, response=response) - return None + return + if output is not None: + await maybe_deferred_to_future( + self.start_itemproc(output, response=response) + ) - def start_itemproc(self, item: Any, *, response: Response | None) -> Deferred[Any]: + @deferred_f_from_coro_f + async def start_itemproc(self, item: Any, *, response: Response | None) -> None: """Send *item* to the item pipelines for processing. *response* is the source of the item data. If the item does not come @@ -313,85 +377,56 @@ class Scraper: assert self.slot is not None # typing assert self.crawler.spider is not None # typing self.slot.itemproc_size += 1 - dfd = self.itemproc.process_item(item, self.crawler.spider) - dfd.addBoth(self._itemproc_finished, item, response, self.crawler.spider) - return dfd - - def _log_download_errors( - self, - spider_failure: Failure, - download_failure: Failure, - request: Request, - spider: Spider, - ) -> Failure | None: - """Log and silence errors that come from the engine (typically download - errors that got propagated thru here). - - spider_failure: the value passed into the errback of self.call_spider() - download_failure: the value passed into _scrape2() from - ExecutionEngine._handle_downloader_output() as "result" - """ - if not download_failure.check(IgnoreRequest): - if download_failure.frames: - logkws = self.logformatter.download_error( - download_failure, request, spider - ) + try: + output = await maybe_deferred_to_future( + self.itemproc.process_item(item, self.crawler.spider) + ) + except DropItem as ex: + logkws = self.logformatter.dropped(item, ex, response, self.crawler.spider) + if logkws is not None: logger.log( - *logformatter_adapter(logkws), - extra={"spider": spider}, - exc_info=failure_to_exc_info(download_failure), + *logformatter_adapter(logkws), extra={"spider": self.crawler.spider} ) - else: - errmsg = download_failure.getErrorMessage() - if errmsg: - logkws = self.logformatter.download_error( - download_failure, request, spider, errmsg - ) - logger.log( - *logformatter_adapter(logkws), - extra={"spider": spider}, - ) - - if spider_failure is not download_failure: - return spider_failure - return None - - def _itemproc_finished( - self, output: Any, item: Any, response: Response | None, spider: Spider - ) -> Deferred[Any]: - """ItemProcessor finished for the given ``item`` and returned ``output``""" - assert self.slot is not None # typing - self.slot.itemproc_size -= 1 - if isinstance(output, Failure): - ex = output.value - if isinstance(ex, DropItem): - logkws = self.logformatter.dropped(item, ex, response, spider) - if logkws is not None: - logger.log(*logformatter_adapter(logkws), extra={"spider": spider}) - return self.signals.send_catch_log_deferred( + await maybe_deferred_to_future( + self.signals.send_catch_log_deferred( signal=signals.item_dropped, item=item, response=response, - spider=spider, - exception=output.value, + spider=self.crawler.spider, + exception=ex, ) - assert ex - logkws = self.logformatter.item_error(item, ex, response, spider) + ) + except Exception as ex: + logkws = self.logformatter.item_error( + item, ex, response, self.crawler.spider + ) logger.log( *logformatter_adapter(logkws), - extra={"spider": spider}, - exc_info=failure_to_exc_info(output), + extra={"spider": self.crawler.spider}, + exc_info=True, ) - return self.signals.send_catch_log_deferred( - signal=signals.item_error, - item=item, - response=response, - spider=spider, - failure=output, + await maybe_deferred_to_future( + self.signals.send_catch_log_deferred( + signal=signals.item_error, + item=item, + response=response, + spider=self.crawler.spider, + failure=Failure(), + ) ) - logkws = self.logformatter.scraped(output, response, spider) - if logkws is not None: - logger.log(*logformatter_adapter(logkws), extra={"spider": spider}) - return self.signals.send_catch_log_deferred( - signal=signals.item_scraped, item=output, response=response, spider=spider - ) + else: + logkws = self.logformatter.scraped(output, response, self.crawler.spider) + if logkws is not None: + logger.log( + *logformatter_adapter(logkws), extra={"spider": self.crawler.spider} + ) + await maybe_deferred_to_future( + self.signals.send_catch_log_deferred( + signal=signals.item_scraped, + item=output, + response=response, + spider=self.crawler.spider, + ) + ) + finally: + self.slot.itemproc_size -= 1 diff --git a/scrapy/core/spidermw.py b/scrapy/core/spidermw.py index 7208bfe67..e9b1d4a1e 100644 --- a/scrapy/core/spidermw.py +++ b/scrapy/core/spidermw.py @@ -41,7 +41,8 @@ logger = logging.getLogger(__name__) _T = TypeVar("_T") ScrapeFunc = Callable[ - [Union[Response, Failure], Request, Spider], Union[Iterable[_T], AsyncIterator[_T]] + [Union[Response, Failure], Request], + Deferred[Union[Iterable[_T], AsyncIterator[_T]]], ] @@ -77,16 +78,15 @@ class SpiderMiddlewareManager(MiddlewareManager): ] if deprecated_middlewares and modern_middlewares: raise ValueError( - f"You are trying to combine spider middlewares that only " - f"define the deprecated process_start_requests() method " - f"({deprecated_middlewares}) with spider middlewares that " - f"only define the process_start() method " - f"({modern_middlewares}). This is not possible. You must " - f"either disable or make universal 1 of those 2 sets of " - f"spider middlewares. Making a spider middleware universal " - f"means having it define both methods. See the release notes " - f"of Scrapy VERSION for details: " - f"https://docs.scrapy.org/en/VERSION/news.html" + "You are trying to combine spider middlewares that only " + "define the deprecated process_start_requests() method () " + "with spider middlewares that only define the " + "process_start() method (). This is not possible. You must " + "either disable or make universal 1 of those 2 sets of " + "spider middlewares. Making a spider middleware universal " + "means having it define both methods. See the release notes " + "of Scrapy VERSION for details: " + "https://docs.scrapy.org/en/VERSION/news.html" ) self._use_start_requests = bool(deprecated_middlewares) @@ -137,7 +137,7 @@ class SpiderMiddlewareManager(MiddlewareManager): response: Response, request: Request, spider: Spider, - ) -> Iterable[_T] | AsyncIterator[_T]: + ) -> Deferred[Iterable[_T] | AsyncIterator[_T]]: for method in self.methods["process_spider_input"]: method = cast(Callable, method) try: @@ -151,8 +151,8 @@ class SpiderMiddlewareManager(MiddlewareManager): except _InvalidOutput: raise except Exception: - return scrape_func(Failure(), request, spider) - return scrape_func(response, request, spider) + return scrape_func(Failure(), request) + return scrape_func(response, request) def _evaluate_iterable( self, @@ -388,20 +388,18 @@ class SpiderMiddlewareManager(MiddlewareManager): dfd2.addErrback(process_spider_exception) return dfd2 - @deferred_f_from_coro_f async def process_start(self, spider: Spider) -> AsyncIterator[Any] | None: try: self._check_deprecated_start_requests_use(spider) except ValueError as exception: logger.error(exception) return None - start: AsyncIterator[Any] if self._use_start_requests: sync_start = iter(spider.start_requests()) sync_start = await maybe_deferred_to_future( self._process_chain("process_start_requests", sync_start, spider) ) - start = as_async_generator(sync_start) + start: AsyncIterator[Any] = as_async_generator(sync_start) else: error_found = False for fn in (spider.start, *self.methods["process_start"]): diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index ac754819c..34d1262ea 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -160,8 +160,16 @@ class ScrapyPriorityQueue: downstream_queue_cls: type[QueueProtocol], key: str, startprios: Iterable[int] = (), + *, + start_queue_cls: type[QueueProtocol] | None = None, ) -> Self: - return cls(crawler, downstream_queue_cls, key, startprios) + return cls( + crawler, + downstream_queue_cls, + key, + startprios, + start_queue_cls=start_queue_cls, + ) def __init__( self, @@ -169,11 +177,15 @@ class ScrapyPriorityQueue: downstream_queue_cls: type[QueueProtocol], key: str, startprios: Iterable[int] = (), + *, + start_queue_cls: type[QueueProtocol] | None = None, ): self.crawler: Crawler = crawler self.downstream_queue_cls: type[QueueProtocol] = downstream_queue_cls + self._start_queue_cls: type[QueueProtocol] | None = start_queue_cls self.key: str = key self.queues: dict[int, QueueProtocol] = {} + self._start_queues: dict[int, QueueProtocol] = {} self.curprio: int | None = None self.init_prios(startprios) @@ -182,7 +194,13 @@ class ScrapyPriorityQueue: return for priority in startprios: - self.queues[priority] = self.qfactory(priority) + q = self.qfactory(priority) + if q: + self.queues[priority] = q + if self._start_queue_cls: + q = self._sqfactory(priority) + if q: + self._start_queues[priority] = q self.curprio = min(startprios) @@ -193,29 +211,66 @@ class ScrapyPriorityQueue: self.key + "/" + str(key), ) + def _sqfactory(self, key: int) -> QueueProtocol: + assert self._start_queue_cls is not None + return build_from_crawler( + self._start_queue_cls, + self.crawler, + f"{self.key}/{key}s", + ) + def priority(self, request: Request) -> int: return -request.priority def push(self, request: Request) -> None: priority = self.priority(request) - if priority not in self.queues: - self.queues[priority] = self.qfactory(priority) - q = self.queues[priority] + is_start_request = request.meta.get("is_start_request", False) + if is_start_request and self._start_queue_cls: + if priority not in self._start_queues: + self._start_queues[priority] = self._sqfactory(priority) + q = self._start_queues[priority] + else: + if priority not in self.queues: + self.queues[priority] = self.qfactory(priority) + q = self.queues[priority] q.push(request) # this may fail (eg. serialization error) if self.curprio is None or priority < self.curprio: self.curprio = priority def pop(self) -> Request | None: - if self.curprio is None: - return None - q = self.queues[self.curprio] - m = q.pop() - if not q: - del self.queues[self.curprio] - q.close() - prios = [p for p, q in self.queues.items() if q] - self.curprio = min(prios) if prios else None - return m + while self.curprio is not None: + if self._start_queues: + try: + q = self._start_queues[self.curprio] + except KeyError: + pass + else: + m = q.pop() + if not q: + del self._start_queues[self.curprio] + q.close() + return m + try: + q = self.queues[self.curprio] + except KeyError: + self._update_curprio() + else: + m = q.pop() + if not q: + del self.queues[self.curprio] + q.close() + self._update_curprio() + return m + return None + + def _update_curprio(self) -> None: + prios = { + p + for queues in (self.queues, self._start_queues) + for p, q in queues.items() + if q + } + self.curprio = min(prios) if prios else None def peek(self) -> Request | None: """Returns the next object to be returned by :meth:`pop`, @@ -226,19 +281,31 @@ class ScrapyPriorityQueue: """ if self.curprio is None: return None - queue = self.queues[self.curprio] + try: + queue = self._start_queues[self.curprio] + except KeyError: + queue = self.queues[self.curprio] # Protocols can't declare optional members return cast(Request, queue.peek()) # type: ignore[attr-defined] def close(self) -> list[int]: - active: list[int] = [] - for p, q in self.queues.items(): - active.append(p) - q.close() - return active + active: set[int] = set() + for queues in (self.queues, self._start_queues): + for p, q in queues.items(): + active.add(p) + q.close() + return list(active) def __len__(self) -> int: - return sum(len(x) for x in self.queues.values()) if self.queues else 0 + return ( + sum( + len(x) + for queues in (self.queues, self._start_queues) + for x in queues.values() + ) + if self.queues or self._start_queues + else 0 + ) class DownloaderInterface: @@ -281,8 +348,16 @@ class DownloaderAwarePriorityQueue: downstream_queue_cls: type[QueueProtocol], key: str, startprios: dict[str, Iterable[int]] | None = None, + *, + start_queue_cls: type[QueueProtocol] | None = None, ) -> Self: - return cls(crawler, downstream_queue_cls, key, startprios) + return cls( + crawler, + downstream_queue_cls, + key, + startprios, + start_queue_cls=start_queue_cls, + ) def __init__( self, @@ -290,6 +365,8 @@ class DownloaderAwarePriorityQueue: downstream_queue_cls: type[QueueProtocol], key: str, slot_startprios: dict[str, Iterable[int]] | None = None, + *, + start_queue_cls: type[QueueProtocol] | None = None, ): if crawler.settings.getint("CONCURRENT_REQUESTS_PER_IP") != 0: raise ValueError( @@ -309,6 +386,7 @@ class DownloaderAwarePriorityQueue: self._downloader_interface: DownloaderInterface = DownloaderInterface(crawler) self.downstream_queue_cls: type[QueueProtocol] = downstream_queue_cls + self._start_queue_cls: type[QueueProtocol] | None = start_queue_cls self.key: str = key self.crawler: Crawler = crawler @@ -324,6 +402,7 @@ class DownloaderAwarePriorityQueue: self.downstream_queue_cls, self.key + "/" + _path_safe(slot), startprios, + start_queue_cls=self._start_queue_cls, ) def pop(self) -> Request | None: diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 680fded7a..01443fa17 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -305,6 +305,8 @@ SCHEDULER = "scrapy.core.scheduler.Scheduler" SCHEDULER_DISK_QUEUE = "scrapy.squeues.PickleLifoDiskQueue" SCHEDULER_MEMORY_QUEUE = "scrapy.squeues.LifoMemoryQueue" SCHEDULER_PRIORITY_QUEUE = "scrapy.pqueues.ScrapyPriorityQueue" +SCHEDULER_START_DISK_QUEUE = "scrapy.squeues.PickleFifoDiskQueue" +SCHEDULER_START_MEMORY_QUEUE = "scrapy.squeues.FifoMemoryQueue" SCRAPER_SLOT_MAX_ACTIVE_SIZE = 5000000 @@ -315,6 +317,7 @@ SPIDER_MIDDLEWARES = {} SPIDER_MIDDLEWARES_BASE = { # Engine side + "scrapy.spidermiddlewares.start.StartSpiderMiddleware": 25, "scrapy.spidermiddlewares.httperror.HttpErrorMiddleware": 50, "scrapy.spidermiddlewares.referer.RefererMiddleware": 700, "scrapy.spidermiddlewares.urllength.UrlLengthMiddleware": 800, diff --git a/scrapy/spidermiddlewares/base.py b/scrapy/spidermiddlewares/base.py new file mode 100644 index 000000000..cfb50c599 --- /dev/null +++ b/scrapy/spidermiddlewares/base.py @@ -0,0 +1,108 @@ +from __future__ import annotations + +from typing import TYPE_CHECKING, Any + +from scrapy import Request, Spider + +if TYPE_CHECKING: + from collections.abc import AsyncIterator, Iterable + + # typing.Self requires Python 3.11 + from typing_extensions import Self + + from scrapy.crawler import Crawler + from scrapy.http import Response + + +class BaseSpiderMiddleware: + """Optional base class for spider middlewares. + + This class provides helper methods for asynchronous + ``process_spider_output()`` and ``process_start()`` methods. Middlewares + that don't have either of these methods don't need to use this class. + + You can override the + :meth:`~scrapy.spidermiddlewares.base.BaseSpiderMiddleware.get_processed_request` + method to add processing code for requests and the + :meth:`~scrapy.spidermiddlewares.base.BaseSpiderMiddleware.get_processed_item` + method to add processing code for items. These methods take a single + request or item from the spider output iterable and return a request or + item (the same or a new one), or ``None`` to remove this request or item + from the processing. + """ + + def __init__(self, crawler: Crawler): + self.crawler: Crawler = crawler + + @classmethod + def from_crawler(cls, crawler: Crawler) -> Self: + return cls(crawler) + + def process_start_requests( + self, start: Iterable[Any], spider: Spider + ) -> Iterable[Any]: + for o in start: + if (o := self._get_processed(o, None)) is not None: + yield o + + async def process_start(self, start: AsyncIterator[Any]) -> AsyncIterator[Any]: + async for o in start: + if (o := self._get_processed(o, None)) is not None: + yield o + + def process_spider_output( + self, response: Response, result: Iterable[Any], spider: Spider + ) -> Iterable[Any]: + for o in result: + if (o := self._get_processed(o, response)) is not None: + yield o + + async def process_spider_output_async( + self, response: Response, result: AsyncIterator[Any], spider: Spider + ) -> AsyncIterator[Any]: + async for o in result: + if (o := self._get_processed(o, response)) is not None: + yield o + + def _get_processed(self, o: Any, response: Response | None) -> Any: + if isinstance(o, Request): + return self.get_processed_request(o, response) + return self.get_processed_item(o, response) + + def get_processed_request( + self, request: Request, response: Response | None + ) -> Request | None: + """Return a processed request from the spider output. + + This method is called with a single request from the start seeds or the + spider output. It should return the same or a different request, or + ``None`` to ignore it. + + :param request: the input request + :type request: :class:`~scrapy.Request` object + + :param response: the response being processed + :type response: :class:`~scrapy.http.Response` object or ``None`` for + start seeds + + :return: the processed request or ``None`` + """ + return request + + def get_processed_item(self, item: Any, response: Response | None) -> Any: + """Return a processed item from the spider output. + + This method is called with a single item from the start seeds or the + spider output. It should return the same or a different item, or + ``None`` to ignore it. + + :param item: the input item + :type item: item object + + :param response: the response being processed + :type response: :class:`~scrapy.http.Response` object or ``None`` for + start seeds + + :return: the processed item or ``None`` + """ + return item diff --git a/scrapy/spidermiddlewares/depth.py b/scrapy/spidermiddlewares/depth.py index 5760411c0..6b115ebe6 100644 --- a/scrapy/spidermiddlewares/depth.py +++ b/scrapy/spidermiddlewares/depth.py @@ -9,7 +9,7 @@ from __future__ import annotations import logging from typing import TYPE_CHECKING, Any -from scrapy.http import Request, Response +from scrapy.spidermiddlewares.base import BaseSpiderMiddleware if TYPE_CHECKING: from collections.abc import AsyncIterator, Iterable @@ -19,14 +19,17 @@ if TYPE_CHECKING: from scrapy import Spider from scrapy.crawler import Crawler + from scrapy.http import Request, Response from scrapy.statscollectors import StatsCollector logger = logging.getLogger(__name__) -class DepthMiddleware: - def __init__( +class DepthMiddleware(BaseSpiderMiddleware): + crawler: Crawler + + def __init__( # pylint: disable=super-init-not-called self, maxdepth: int, stats: StatsCollector, @@ -45,21 +48,22 @@ class DepthMiddleware: verbose = settings.getbool("DEPTH_STATS_VERBOSE") prio = settings.getint("DEPTH_PRIORITY") assert crawler.stats - return cls(maxdepth, crawler.stats, verbose, prio) + o = cls(maxdepth, crawler.stats, verbose, prio) + o.crawler = crawler + return o def process_spider_output( self, response: Response, result: Iterable[Any], spider: Spider ) -> Iterable[Any]: self._init_depth(response, spider) - return (r for r in result if self._filter(r, response, spider)) + yield from super().process_spider_output(response, result, spider) async def process_spider_output_async( self, response: Response, result: AsyncIterator[Any], spider: Spider ) -> AsyncIterator[Any]: self._init_depth(response, spider) - async for r in result: - if self._filter(r, response, spider): - yield r + async for o in super().process_spider_output_async(response, result, spider): + yield o def _init_depth(self, response: Response, spider: Spider) -> None: # base case (depth=0) @@ -68,9 +72,12 @@ class DepthMiddleware: if self.verbose_stats: self.stats.inc_value("request_depth_count/0", spider=spider) - def _filter(self, request: Any, response: Response, spider: Spider) -> bool: - if not isinstance(request, Request): - return True + def get_processed_request( + self, request: Request, response: Response | None + ) -> Request | None: + if response is None: + # start requests + return request depth = response.meta["depth"] + 1 request.meta["depth"] = depth if self.prio: @@ -79,10 +86,12 @@ class DepthMiddleware: logger.debug( "Ignoring link (depth > %(maxdepth)d): %(requrl)s ", {"maxdepth": self.maxdepth, "requrl": request.url}, - extra={"spider": spider}, + extra={"spider": self.crawler.spider}, ) - return False + return None if self.verbose_stats: - self.stats.inc_value(f"request_depth_count/{depth}", spider=spider) - self.stats.max_value("request_depth_max", depth, spider=spider) - return True + self.stats.inc_value( + f"request_depth_count/{depth}", spider=self.crawler.spider + ) + self.stats.max_value("request_depth_max", depth, spider=self.crawler.spider) + return request diff --git a/scrapy/spidermiddlewares/offsite.py b/scrapy/spidermiddlewares/offsite.py index 060be1975..2463275d5 100644 --- a/scrapy/spidermiddlewares/offsite.py +++ b/scrapy/spidermiddlewares/offsite.py @@ -9,11 +9,11 @@ from __future__ import annotations import logging import re import warnings -from typing import TYPE_CHECKING, Any +from typing import TYPE_CHECKING from scrapy import Spider, signals from scrapy.exceptions import ScrapyDeprecationWarning -from scrapy.http import Request, Response +from scrapy.spidermiddlewares.base import BaseSpiderMiddleware from scrapy.utils.httpobj import urlparse_cached warnings.warn( @@ -23,61 +23,55 @@ warnings.warn( ) if TYPE_CHECKING: - from collections.abc import AsyncIterator, Iterable - # typing.Self requires Python 3.11 from typing_extensions import Self from scrapy.crawler import Crawler + from scrapy.http import Request, Response from scrapy.statscollectors import StatsCollector logger = logging.getLogger(__name__) -class OffsiteMiddleware: - def __init__(self, stats: StatsCollector): +class OffsiteMiddleware(BaseSpiderMiddleware): + crawler: Crawler + + def __init__(self, stats: StatsCollector): # pylint: disable=super-init-not-called self.stats: StatsCollector = stats @classmethod def from_crawler(cls, crawler: Crawler) -> Self: assert crawler.stats o = cls(crawler.stats) + o.crawler = crawler crawler.signals.connect(o.spider_opened, signal=signals.spider_opened) return o - def process_spider_output( - self, response: Response, result: Iterable[Any], spider: Spider - ) -> Iterable[Any]: - return (r for r in result if self._filter(r, spider)) - - async def process_spider_output_async( - self, response: Response, result: AsyncIterator[Any], spider: Spider - ) -> AsyncIterator[Any]: - async for r in result: - if self._filter(r, spider): - yield r - - def _filter(self, request: Any, spider: Spider) -> bool: - if not isinstance(request, Request): - return True + def get_processed_request( + self, request: Request, response: Response | None + ) -> Request | None: + if response is None: + # skip start requests for backward compatibility + return request + assert self.crawler.spider if ( request.dont_filter or request.meta.get("allow_offsite") - or self.should_follow(request, spider) + or self.should_follow(request, self.crawler.spider) ): - return True + return request domain = urlparse_cached(request).hostname if domain and domain not in self.domains_seen: self.domains_seen.add(domain) logger.debug( "Filtered offsite request to %(domain)r: %(request)s", {"domain": domain, "request": request}, - extra={"spider": spider}, + extra={"spider": self.crawler.spider}, ) - self.stats.inc_value("offsite/domains", spider=spider) - self.stats.inc_value("offsite/filtered", spider=spider) - return False + self.stats.inc_value("offsite/domains", spider=self.crawler.spider) + self.stats.inc_value("offsite/filtered", spider=self.crawler.spider) + return None def should_follow(self, request: Request, spider: Spider) -> bool: regex = self.host_regex diff --git a/scrapy/spidermiddlewares/referer.py b/scrapy/spidermiddlewares/referer.py index c824579b0..f5d406c13 100644 --- a/scrapy/spidermiddlewares/referer.py +++ b/scrapy/spidermiddlewares/referer.py @@ -6,7 +6,7 @@ originated it. from __future__ import annotations import warnings -from typing import TYPE_CHECKING, Any, cast +from typing import TYPE_CHECKING, cast from urllib.parse import urlparse from w3lib.url import safe_url_string @@ -14,13 +14,12 @@ from w3lib.url import safe_url_string from scrapy import Spider, signals from scrapy.exceptions import NotConfigured from scrapy.http import Request, Response +from scrapy.spidermiddlewares.base import BaseSpiderMiddleware from scrapy.utils.misc import load_object from scrapy.utils.python import to_unicode from scrapy.utils.url import strip_url if TYPE_CHECKING: - from collections.abc import AsyncIterator, Iterable - # typing.Self requires Python 3.11 from typing_extensions import Self @@ -327,8 +326,8 @@ def _load_policy_class( return None -class RefererMiddleware: - def __init__(self, settings: BaseSettings | None = None): +class RefererMiddleware(BaseSpiderMiddleware): + def __init__(self, settings: BaseSettings | None = None): # pylint: disable=super-init-not-called self.default_policy: type[ReferrerPolicy] = DefaultReferrerPolicy if settings is not None: settings_policy = _load_policy_class(settings.get("REFERRER_POLICY")) @@ -370,23 +369,16 @@ class RefererMiddleware: cls = _load_policy_class(policy_name, warning_only=True) return cls() if cls else self.default_policy() - def process_spider_output( - self, response: Response, result: Iterable[Any], spider: Spider - ) -> Iterable[Any]: - return (self._set_referer(r, response) for r in result) - - async def process_spider_output_async( - self, response: Response, result: AsyncIterator[Any], spider: Spider - ) -> AsyncIterator[Any]: - async for r in result: - yield self._set_referer(r, response) - - def _set_referer(self, r: Any, response: Response) -> Any: - if isinstance(r, Request): - referrer = self.policy(response, r).referrer(response.url, r.url) - if referrer is not None: - r.headers.setdefault("Referer", referrer) - return r + def get_processed_request( + self, request: Request, response: Response | None + ) -> Request | None: + if response is None: + # start requests + return request + referrer = self.policy(response, request).referrer(response.url, request.url) + if referrer is not None: + request.headers.setdefault("Referer", referrer) + return request def request_scheduled(self, request: Request, spider: Spider) -> None: # check redirected request to patch "Referer" header if necessary diff --git a/scrapy/spidermiddlewares/start.py b/scrapy/spidermiddlewares/start.py new file mode 100644 index 000000000..5d76b60d2 --- /dev/null +++ b/scrapy/spidermiddlewares/start.py @@ -0,0 +1,31 @@ +from __future__ import annotations + +from typing import TYPE_CHECKING + +from .base import BaseSpiderMiddleware + +if TYPE_CHECKING: + from scrapy.http import Request + from scrapy.http.response import Response + + +class StartSpiderMiddleware(BaseSpiderMiddleware): + """Set :reqmeta:`is_start_request`. + + .. reqmeta:: is_start_request + + is_start_request + ---------------- + + :attr:`~scrapy.Request.meta` key that is set to ``True`` in :ref:`start + requests `, allowing you to tell start requests apart from + other requests, e.g. in :ref:`downloader middlewares + `. + """ + + def get_processed_request( + self, request: Request, response: Response | None + ) -> Request | None: + if response is None: + request.meta.setdefault("is_start_request", True) + return request diff --git a/scrapy/spidermiddlewares/urllength.py b/scrapy/spidermiddlewares/urllength.py index 9259fea3c..5590165a5 100644 --- a/scrapy/spidermiddlewares/urllength.py +++ b/scrapy/spidermiddlewares/urllength.py @@ -7,72 +7,49 @@ See documentation in docs/topics/spider-middleware.rst from __future__ import annotations import logging -import warnings -from typing import TYPE_CHECKING, Any +from typing import TYPE_CHECKING -from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning -from scrapy.http import Request, Response +from scrapy.exceptions import NotConfigured +from scrapy.spidermiddlewares.base import BaseSpiderMiddleware if TYPE_CHECKING: - from collections.abc import AsyncIterator, Iterable - # typing.Self requires Python 3.11 from typing_extensions import Self - from scrapy import Spider from scrapy.crawler import Crawler - from scrapy.settings import BaseSettings + from scrapy.http import Request, Response logger = logging.getLogger(__name__) -class UrlLengthMiddleware: - def __init__(self, maxlength: int): +class UrlLengthMiddleware(BaseSpiderMiddleware): + crawler: Crawler + + def __init__(self, maxlength: int): # pylint: disable=super-init-not-called self.maxlength: int = maxlength - @classmethod - def from_settings(cls, settings: BaseSettings) -> Self: - warnings.warn( - f"{cls.__name__}.from_settings() is deprecated, use from_crawler() instead.", - category=ScrapyDeprecationWarning, - stacklevel=2, - ) - return cls._from_settings(settings) - @classmethod def from_crawler(cls, crawler: Crawler) -> Self: - return cls._from_settings(crawler.settings) - - @classmethod - def _from_settings(cls, settings: BaseSettings) -> Self: - maxlength = settings.getint("URLLENGTH_LIMIT") + maxlength = crawler.settings.getint("URLLENGTH_LIMIT") if not maxlength: raise NotConfigured - return cls(maxlength) + o = cls(maxlength) + o.crawler = crawler + return o - def process_spider_output( - self, response: Response, result: Iterable[Any], spider: Spider - ) -> Iterable[Any]: - return (r for r in result if self._filter(r, spider)) - - async def process_spider_output_async( - self, response: Response, result: AsyncIterator[Any], spider: Spider - ) -> AsyncIterator[Any]: - async for r in result: - if self._filter(r, spider): - yield r - - def _filter(self, request: Any, spider: Spider) -> bool: - if isinstance(request, Request) and len(request.url) > self.maxlength: - logger.info( - "Ignoring link (url length > %(maxlength)d): %(url)s ", - {"maxlength": self.maxlength, "url": request.url}, - extra={"spider": spider}, - ) - assert spider.crawler.stats - spider.crawler.stats.inc_value( - "urllength/request_ignored_count", spider=spider - ) - return False - return True + def get_processed_request( + self, request: Request, response: Response | None + ) -> Request | None: + if len(request.url) <= self.maxlength: + return request + logger.info( + "Ignoring link (url length > %(maxlength)d): %(url)s ", + {"maxlength": self.maxlength, "url": request.url}, + extra={"spider": self.crawler.spider}, + ) + assert self.crawler.stats + self.crawler.stats.inc_value( + "urllength/request_ignored_count", spider=self.crawler.spider + ) + return None diff --git a/scrapy/squeues.py b/scrapy/squeues.py index 6d73d66f0..a041f36ba 100644 --- a/scrapy/squeues.py +++ b/scrapy/squeues.py @@ -177,8 +177,6 @@ PickleFifoDiskQueue = _scrapy_serialization_queue(_PickleFifoSerializationDiskQu #: LIFO_ disk queue that serializes :ref:`requests ` using #: :mod:`pickle`. -#: -#: .. _LIFO: https://en.wikipedia.org/wiki/LIFO_(computing) PickleLifoDiskQueue = _scrapy_serialization_queue(_PickleLifoSerializationDiskQueue) #: FIFO_ disk queue that serializes :ref:`requests ` using diff --git a/scrapy/utils/defer.py b/scrapy/utils/defer.py index c39837716..6e1687f3e 100644 --- a/scrapy/utils/defer.py +++ b/scrapy/utils/defer.py @@ -14,7 +14,11 @@ from types import CoroutineType from typing import TYPE_CHECKING, Any, Generic, TypeVar, Union, cast, overload from twisted.internet import defer -from twisted.internet.defer import Deferred, DeferredList, ensureDeferred +from twisted.internet.defer import ( + Deferred, + DeferredList, + ensureDeferred, +) from twisted.internet.task import Cooperator from twisted.python import failure @@ -36,6 +40,9 @@ _T = TypeVar("_T") _T2 = TypeVar("_T2") +_DEFER_DELAY = 0.1 + + def defer_fail(_failure: Failure) -> Deferred[Any]: """Same as twisted.internet.defer.fail but delay calling errback until next reactor loop @@ -46,7 +53,7 @@ def defer_fail(_failure: Failure) -> Deferred[Any]: from twisted.internet import reactor d: Deferred[Any] = Deferred() - reactor.callLater(0.1, d.errback, _failure) + reactor.callLater(_DEFER_DELAY, d.errback, _failure) return d @@ -60,7 +67,16 @@ def defer_succeed(result: _T) -> Deferred[_T]: from twisted.internet import reactor d: Deferred[_T] = Deferred() - reactor.callLater(0.1, d.callback, result) + reactor.callLater(_DEFER_DELAY, d.callback, result) + return d + + +def _defer_sleep() -> Deferred[None]: + """Like ``defer_succeed`` and ``defer_fail`` but doesn't call any real callbacks.""" + from twisted.internet import reactor + + d: Deferred[None] = Deferred() + reactor.callLater(_DEFER_DELAY, d.callback, None) return d @@ -338,7 +354,7 @@ async def aiter_errback( **kw: _P.kwargs, ) -> AsyncIterator[_T]: """Wraps an async iterable calling an errback if an error is caught while - iterating it. Similar to scrapy.utils.defer.iter_errback() + iterating it. Similar to :func:`scrapy.utils.defer.iter_errback`. """ it = aiterable.__aiter__() while True: diff --git a/tests/CrawlerProcess/sleeping.py b/tests/CrawlerProcess/sleeping.py index 45479ea4f..cb8f869e1 100644 --- a/tests/CrawlerProcess/sleeping.py +++ b/tests/CrawlerProcess/sleeping.py @@ -1,3 +1,5 @@ +import sys + from twisted.internet.defer import Deferred import scrapy @@ -14,7 +16,7 @@ class SleepingSpider(scrapy.Spider): from twisted.internet import reactor d = Deferred() - reactor.callLater(int(self.sleep), d.callback, None) + reactor.callLater(int(sys.argv[1]), d.callback, None) await maybe_deferred_to_future(d) diff --git a/tests/test_crawl.py b/tests/test_crawl.py index 7793505da..b90706027 100644 --- a/tests/test_crawl.py +++ b/tests/test_crawl.py @@ -1,3 +1,5 @@ +from __future__ import annotations + import json import logging import unittest @@ -14,7 +16,7 @@ from twisted.trial.unittest import TestCase from scrapy import signals from scrapy.crawler import CrawlerRunner -from scrapy.exceptions import StopDownload +from scrapy.exceptions import CloseSpider, StopDownload from scrapy.http import Request from scrapy.http.response import Response from scrapy.utils.python import to_unicode @@ -186,11 +188,18 @@ class TestCrawl(TestCase): @defer.inlineCallbacks def test_start_items(self): + items = [] + + def _on_item_scraped(item): + items.append(item) + with LogCapture("scrapy", level=logging.ERROR) as log: crawler = get_crawler(StartItemSpider) + crawler.signals.connect(_on_item_scraped, signals.item_scraped) yield crawler.crawl(mockserver=self.mockserver) assert len(log.records) == 0 + assert items == [{"name": "test item"}] @defer.inlineCallbacks def test_start_unsupported_output(self): @@ -198,11 +207,20 @@ class TestCrawl(TestCase): potentially expensive call to itemadapter.is_item(), and letting instead things fail when ItemAdapter is actually used on the corresponding non-item object.""" + + items = [] + + def _on_item_scraped(item): + items.append(item) + with LogCapture("scrapy", level=logging.ERROR) as log: crawler = get_crawler(StartGoodAndBadOutput) + crawler.signals.connect(_on_item_scraped, signals.item_scraped) yield crawler.crawl(mockserver=self.mockserver) assert len(log.records) == 0 + assert len(items) == 3 + assert not any(isinstance(item, Request) for item in items) @defer.inlineCallbacks def test_start_dupes(self): @@ -692,3 +710,100 @@ class TestCrawlSpider(TestCase): assert crawler.spider.meta[ "failure" ].value.response.headers == crawler.spider.meta.get("headers_received") + + @defer.inlineCallbacks + def test_spider_errback(self): + failures = [] + + def eb(failure: Failure) -> Failure: + failures.append(failure) + return failure + + crawler = get_crawler(SingleRequestSpider) + with LogCapture() as log: + yield crawler.crawl( + seed=self.mockserver.url("/status?n=400"), errback_func=eb + ) + assert len(failures) == 1 + assert "HTTP status code is not handled or not allowed" in str(log) + assert "Spider error processing" not in str(log) + + @defer.inlineCallbacks + def test_spider_errback_silence(self): + failures = [] + + def eb(failure: Failure) -> None: + failures.append(failure) + + crawler = get_crawler(SingleRequestSpider) + with LogCapture() as log: + yield crawler.crawl( + seed=self.mockserver.url("/status?n=400"), errback_func=eb + ) + assert len(failures) == 1 + assert "HTTP status code is not handled or not allowed" not in str(log) + assert "Spider error processing" not in str(log) + + @defer.inlineCallbacks + def test_spider_errback_exception(self): + def eb(failure: Failure) -> None: + raise ValueError("foo") + + crawler = get_crawler(SingleRequestSpider) + with LogCapture() as log: + yield crawler.crawl( + seed=self.mockserver.url("/status?n=400"), errback_func=eb + ) + assert "Spider error processing" in str(log) + + @defer.inlineCallbacks + def test_spider_errback_downloader_error(self): + failures = [] + + def eb(failure: Failure) -> Failure: + failures.append(failure) + return failure + + crawler = get_crawler(SingleRequestSpider) + with LogCapture() as log: + yield crawler.crawl( + seed=self.mockserver.url("/drop?abort=1"), errback_func=eb + ) + assert len(failures) == 1 + assert "Error downloading" in str(log) + assert "Spider error processing" not in str(log) + + @defer.inlineCallbacks + def test_spider_errback_exception_downloader_error(self): + def eb(failure: Failure) -> None: + raise ValueError("foo") + + crawler = get_crawler(SingleRequestSpider) + with LogCapture() as log: + yield crawler.crawl( + seed=self.mockserver.url("/drop?abort=1"), errback_func=eb + ) + assert "Error downloading" in str(log) + assert "Spider error processing" in str(log) + + @defer.inlineCallbacks + def test_raise_closespider(self): + def cb(response): + raise CloseSpider + + crawler = get_crawler(SingleRequestSpider) + with LogCapture() as log: + yield crawler.crawl(seed=self.mockserver.url("/"), callback_func=cb) + assert "Closing spider (cancelled)" in str(log) + assert "Spider error processing" not in str(log) + + @defer.inlineCallbacks + def test_raise_closespider_reason(self): + def cb(response): + raise CloseSpider("my_reason") + + crawler = get_crawler(SingleRequestSpider) + with LogCapture() as log: + yield crawler.crawl(seed=self.mockserver.url("/"), callback_func=cb) + assert "Closing spider (my_reason)" in str(log) + assert "Spider error processing" not in str(log) diff --git a/tests/test_crawler.py b/tests/test_crawler.py index f2eb5efcb..4c179c65b 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -891,7 +891,7 @@ class TestCrawlerProcessSubprocess(ScriptRunnerMixin, unittest.TestCase): def test_shutdown_graceful(self): sig = signal.SIGINT if sys.platform != "win32" else signal.SIGBREAK - args = self.get_script_args("sleeping.py", "-a", "sleep=3") + args = self.get_script_args("sleeping.py", "3") p = PopenSpawn(args, timeout=5) p.expect_exact("Spider opened") p.expect_exact("Crawled (200)") @@ -905,7 +905,7 @@ class TestCrawlerProcessSubprocess(ScriptRunnerMixin, unittest.TestCase): from twisted.internet import reactor sig = signal.SIGINT if sys.platform != "win32" else signal.SIGBREAK - args = self.get_script_args("sleeping.py", "-a", "sleep=10") + args = self.get_script_args("sleeping.py", "10") p = PopenSpawn(args, timeout=5) p.expect_exact("Spider opened") p.expect_exact("Crawled (200)") diff --git a/tests/test_downloadermiddleware.py b/tests/test_downloadermiddleware.py index 6c061b330..8ae160f8a 100644 --- a/tests/test_downloadermiddleware.py +++ b/tests/test_downloadermiddleware.py @@ -1,22 +1,18 @@ +from __future__ import annotations + import asyncio from gzip import BadGzipFile from unittest import mock import pytest -from twisted.internet import defer -from twisted.internet.defer import Deferred -from twisted.python.failure import Failure +from twisted.internet.defer import Deferred, succeed from twisted.trial.unittest import TestCase from scrapy.core.downloader.middleware import DownloaderMiddlewareManager from scrapy.exceptions import _InvalidOutput from scrapy.http import Request, Response from scrapy.spiders import Spider -from scrapy.utils.defer import ( - deferred_f_from_coro_f, - deferred_to_future, - maybe_deferred_to_future, -) +from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future from scrapy.utils.python import to_bytes from scrapy.utils.test import get_crawler, get_from_asyncio_queue @@ -34,26 +30,22 @@ class TestManagerBase(TestCase): def tearDown(self): return self.crawler.engine.close_spider(self.spider) - async def _download(self, request, response=None): + async def _download( + self, request: Request, response: Response | None = None + ) -> Response | Request: """Executes downloader mw manager's download method and returns - the result (Request or Response) or raise exception in case of + the result (Request or Response) or raises exception in case of failure. """ if not response: response = Response(request.url) - def download_func(request, spider): - return response + def download_func(request: Request, spider: Spider) -> Deferred[Response]: + return succeed(response) - dfd = self.mwman.download(download_func, request, self.spider) - # catch deferred result and return the value - results = [] - dfd.addBoth(results.append) - await maybe_deferred_to_future(dfd) - ret = results[0] - if isinstance(ret, Failure): - ret.raiseException() - return ret + return await maybe_deferred_to_future( + self.mwman.download(download_func, request, self.spider) + ) class TestDefaults(TestManagerBase): @@ -92,7 +84,7 @@ class TestDefaults(TestManagerBase): "Location": "http://example.com/login", }, ) - ret = await self._download(request=req, response=resp) + ret = await self._download(req, resp) assert isinstance(ret, Request), f"Not redirected: {ret!r}" assert to_bytes(ret.url) == resp.headers["Location"], ( "Not redirected to location header" @@ -114,7 +106,7 @@ class TestDefaults(TestManagerBase): }, ) with pytest.raises(BadGzipFile): - await self._download(request=req, response=resp) + await self._download(req, resp) class TestResponseFromProcessRequest(TestManagerBase): @@ -132,19 +124,17 @@ class TestResponseFromProcessRequest(TestManagerBase): req = Request("http://example.com/index.html") download_func = mock.MagicMock() - dfd = self.mwman.download(download_func, req, self.spider) - results = [] - dfd.addBoth(results.append) - await maybe_deferred_to_future(dfd) - - assert results[0] is resp + result = await maybe_deferred_to_future( + self.mwman.download(download_func, req, self.spider) + ) + assert result is resp assert not download_func.called -class TestProcessRequestInvalidOutput(TestManagerBase): - """Invalid return value for process_request method should raise an exception""" - - def test_invalid_process_request(self): +class TestInvalidOutput(TestManagerBase): + @deferred_f_from_coro_f + async def test_invalid_process_request(self): + """Invalid return value for process_request method should raise an exception""" req = Request("http://example.com/index.html") class InvalidProcessRequestMiddleware: @@ -152,18 +142,12 @@ class TestProcessRequestInvalidOutput(TestManagerBase): return 1 self.mwman._add_middleware(InvalidProcessRequestMiddleware()) - download_func = mock.MagicMock() - dfd = self.mwman.download(download_func, req, self.spider) - results = [] - dfd.addBoth(results.append) - assert isinstance(results[0], Failure) - assert isinstance(results[0].value, _InvalidOutput) + with pytest.raises(_InvalidOutput): + await self._download(req) - -class TestProcessResponseInvalidOutput(TestManagerBase): - """Invalid return value for process_response method should raise an exception""" - - def test_invalid_process_response(self): + @deferred_f_from_coro_f + async def test_invalid_process_response(self): + """Invalid return value for process_response method should raise an exception""" req = Request("http://example.com/index.html") class InvalidProcessResponseMiddleware: @@ -171,18 +155,12 @@ class TestProcessResponseInvalidOutput(TestManagerBase): return 1 self.mwman._add_middleware(InvalidProcessResponseMiddleware()) - download_func = mock.MagicMock() - dfd = self.mwman.download(download_func, req, self.spider) - results = [] - dfd.addBoth(results.append) - assert isinstance(results[0], Failure) - assert isinstance(results[0].value, _InvalidOutput) + with pytest.raises(_InvalidOutput): + await self._download(req) - -class TestProcessExceptionInvalidOutput(TestManagerBase): - """Invalid return value for process_exception method should raise an exception""" - - def test_invalid_process_exception(self): + @deferred_f_from_coro_f + async def test_invalid_process_exception(self): + """Invalid return value for process_exception method should raise an exception""" req = Request("http://example.com/index.html") class InvalidProcessExceptionMiddleware: @@ -193,12 +171,8 @@ class TestProcessExceptionInvalidOutput(TestManagerBase): return 1 self.mwman._add_middleware(InvalidProcessExceptionMiddleware()) - download_func = mock.MagicMock() - dfd = self.mwman.download(download_func, req, self.spider) - results = [] - dfd.addBoth(results.append) - assert isinstance(results[0], Failure) - assert isinstance(results[0].value, _InvalidOutput) + with pytest.raises(_InvalidOutput): + await self._download(req) class TestMiddlewareUsingDeferreds(TestManagerBase): @@ -221,12 +195,10 @@ class TestMiddlewareUsingDeferreds(TestManagerBase): self.mwman._add_middleware(DeferredMiddleware()) req = Request("http://example.com/index.html") download_func = mock.MagicMock() - dfd = self.mwman.download(download_func, req, self.spider) - results = [] - dfd.addBoth(results.append) - await maybe_deferred_to_future(dfd) - - assert results[0] is resp + result = await maybe_deferred_to_future( + self.mwman.download(download_func, req, self.spider) + ) + assert result is resp assert not download_func.called @@ -240,18 +212,16 @@ class TestMiddlewareUsingCoro(TestManagerBase): class CoroMiddleware: async def process_request(self, request, spider): - await defer.succeed(42) + await succeed(42) return resp self.mwman._add_middleware(CoroMiddleware()) req = Request("http://example.com/index.html") download_func = mock.MagicMock() - dfd = self.mwman.download(download_func, req, self.spider) - results = [] - dfd.addBoth(results.append) - await maybe_deferred_to_future(dfd) - - assert results[0] is resp + result = await maybe_deferred_to_future( + self.mwman.download(download_func, req, self.spider) + ) + assert result is resp assert not download_func.called @pytest.mark.only_asyncio @@ -267,10 +237,8 @@ class TestMiddlewareUsingCoro(TestManagerBase): self.mwman._add_middleware(CoroMiddleware()) req = Request("http://example.com/index.html") download_func = mock.MagicMock() - dfd = self.mwman.download(download_func, req, self.spider) - results = [] - dfd.addBoth(results.append) - await deferred_to_future(dfd) - - assert results[0] is resp + result = await maybe_deferred_to_future( + self.mwman.download(download_func, req, self.spider) + ) + assert result is resp assert not download_func.called diff --git a/tests/test_engine.py b/tests/test_engine.py index c8e0d8b25..e85818e68 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -482,8 +482,7 @@ def test_request_scheduled_signal(caplog): if "drop" in request.url: raise IgnoreRequest - spider = MySpider() - crawler = get_crawler(spider.__class__) + crawler = get_crawler(MySpider) engine = ExecutionEngine(crawler, lambda _: None) engine.downloader._slot_gc_loop.stop() scheduler = MemoryScheduler() @@ -497,10 +496,10 @@ def test_request_scheduled_signal(caplog): engine._slot = _Slot(False, Mock()) crawler.signals.connect(signal_handler, request_scheduled) keep_request = Request("https://keep.example") - engine._schedule_request(keep_request, spider) + engine._schedule_request(keep_request) drop_request = Request("https://drop.example") caplog.set_level(DEBUG) - engine._schedule_request(drop_request, spider) + engine._schedule_request(drop_request) assert list(scheduler.queue) == [keep_request], ( f"{list(scheduler.queue)!r} != [{keep_request!r}]" ) diff --git a/tests/test_engine_loop.py b/tests/test_engine_loop.py index 656216e27..2fad6c2fe 100644 --- a/tests/test_engine_loop.py +++ b/tests/test_engine_loop.py @@ -10,6 +10,7 @@ from twisted.trial.unittest import TestCase from scrapy import Request, Spider, signals from scrapy.core.engine import ExecutionEngine +from scrapy.squeues import LifoMemoryQueue from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future from scrapy.utils.test import get_crawler @@ -208,10 +209,10 @@ class RequestSendOrderTestCase(TestCase): def tearDownClass(cls): cls.mockserver.__exit__(None, None, None) # increase if flaky - def request(self, num, response_seconds, download_slots=1): + def request(self, num, response_seconds, download_slots=1, priority=0): url = self.mockserver.url(f"/delay?n={response_seconds}&{num}") meta = {"download_slot": str(num % download_slots)} - return Request(url, meta=meta) + return Request(url, meta=meta, priority=priority) def get_num(self, request_or_response: Request | Response): return int(request_or_response.url.rsplit("&", maxsplit=1)[1]) @@ -264,55 +265,156 @@ class RequestSendOrderTestCase(TestCase): expected_nums = sorted(start_nums + cb_nums) assert actual_nums == expected_nums, f"{actual_nums=} != {expected_nums=}" + @deferred_f_from_coro_f + async def test_default(self): + """By default, start requests take priority over callback requests and + are sent in order. Priority matters, but given the same priority, a + start request takes precedence.""" + nums = [1, 2, 3, 4, 5, 6] + response_seconds = 0 + download_slots = 1 + + def _request(num, priority=0): + return self.request( + num, response_seconds, download_slots, priority=priority + ) + + async def start(spider): + # The first CONCURRENT_REQUESTS start requests are sent + # immediately. + yield _request(1) + + for request in ( + _request(4, priority=1), + _request(6), + ): + spider.crawler.engine.scheduler.enqueue_request(request) + yield _request(5) + yield _request(2, priority=1) + yield _request(3, priority=1) + + def parse(spider, response): + return + yield + + await maybe_deferred_to_future( + self._test_request_order( + start_nums=nums, + settings={"CONCURRENT_REQUESTS": 1}, + response_seconds=response_seconds, + start_fn=start, + parse_fn=parse, + ) + ) + + @deferred_f_from_coro_f + async def test_lifo_start(self): + """Changing the queues of start requests to LIFO, matching the queues + of non-start requests, does not cause all requests to be stored in the + same queue objects, it only affects the order of start requests.""" + nums = [1, 2, 3, 4, 5, 6] + response_seconds = 0 + download_slots = 1 + + def _request(num, priority=0): + return self.request( + num, response_seconds, download_slots, priority=priority + ) + + async def start(spider): + # The first CONCURRENT_REQUESTS start requests are sent + # immediately. + yield _request(1) + + for request in ( + _request(4, priority=1), + _request(6), + ): + spider.crawler.engine.scheduler.enqueue_request(request) + yield _request(5) + yield _request(3, priority=1) + yield _request(2, priority=1) + + def parse(spider, response): + return + yield + + await maybe_deferred_to_future( + self._test_request_order( + start_nums=nums, + settings={ + "CONCURRENT_REQUESTS": 1, + "SCHEDULER_START_MEMORY_QUEUE": "scrapy.squeues.LifoMemoryQueue", + }, + response_seconds=response_seconds, + start_fn=start, + parse_fn=parse, + ) + ) + + @deferred_f_from_coro_f + async def test_shared_queues(self): + """If SCHEDULER_START_*_QUEUE is falsy, start requests and other + requests share the same queue, i.e. start requests are not priorized + over other requests if their priority matches.""" + nums = list(range(1, 14)) + response_seconds = 0 + download_slots = 1 + + def _request(num, priority=0): + return self.request( + num, response_seconds, download_slots, priority=priority + ) + + async def start(spider): + # The first CONCURRENT_REQUESTS start requests are sent + # immediately. + yield _request(1) + + # Below, priority 1 requests are sent first, and requests are sent + # in LIFO order. + + for request in ( + _request(7, priority=1), + _request(6, priority=1), + _request(13), + _request(12), + ): + spider.crawler.engine.scheduler.enqueue_request(request) + + yield _request(11) + yield _request(10) + yield _request(5, priority=1) + yield _request(4, priority=1) + + for request in ( + _request(3, priority=1), + _request(2, priority=1), + _request(9), + _request(8), + ): + spider.crawler.engine.scheduler.enqueue_request(request) + + def parse(spider, response): + return + yield + + await maybe_deferred_to_future( + self._test_request_order( + start_nums=nums, + settings={ + "CONCURRENT_REQUESTS": 1, + "SCHEDULER_START_MEMORY_QUEUE": None, + }, + response_seconds=response_seconds, + start_fn=start, + parse_fn=parse, + ) + ) + # Examples from the “Start requests” section of the documentation about # spiders. - @deferred_f_from_coro_f - async def test_start_requests_first(self): - start_nums = [1, 3, 2] - cb_nums = [4] - response_seconds = self.seconds - download_slots = 1 - - async def start(spider): - for num in start_nums: - request = self.request(num, response_seconds, download_slots) - yield request.replace(priority=1) - - await maybe_deferred_to_future( - self._test_request_order( - start_nums=start_nums, - cb_nums=cb_nums, - settings={"CONCURRENT_REQUESTS": 1}, - response_seconds=response_seconds, - start_fn=start, - ) - ) - - @deferred_f_from_coro_f - async def test_start_requests_first_sorted(self): - start_nums = [1, 2, 3] - cb_nums = [4] - response_seconds = self.seconds - download_slots = 1 - - async def start(spider): - priority = len(start_nums) - for num in start_nums: - request = self.request(num, response_seconds, download_slots) - yield request.replace(priority=priority) - priority -= 1 - - await maybe_deferred_to_future( - self._test_request_order( - start_nums=start_nums, - cb_nums=cb_nums, - settings={"CONCURRENT_REQUESTS": 1}, - response_seconds=response_seconds, - start_fn=start, - ) - ) - @deferred_f_from_coro_f async def test_front_load(self): start_nums = [2, 1] @@ -325,8 +427,8 @@ class RequestSendOrderTestCase(TestCase): assert isinstance(spider.crawler.engine.scheduler, MemoryScheduler) spider.crawler.engine.scheduler.pause() - # By pausing the scheduler, a is scheduled before b is sent, - # and since the scheduler uses a LIFO queue, a is sent first. + # By pausing the scheduler, 1 is scheduled before 2 is sent, + # and provided the scheduler uses a LIFO queue, 1 is sent first. yield self.request(2, response_seconds, download_slots) yield self.request(1, response_seconds, download_slots) spider.crawler.engine.scheduler.unpause() @@ -336,6 +438,9 @@ class RequestSendOrderTestCase(TestCase): start_nums=start_nums, response_seconds=response_seconds, start_fn=start, + settings={ + "SCHEDULER_START_MEMORY_QUEUE": LifoMemoryQueue, + }, ) ) diff --git a/tests/test_spidermiddleware.py b/tests/test_spidermiddleware.py index 9cd7b44ca..db46be7dd 100644 --- a/tests/test_spidermiddleware.py +++ b/tests/test_spidermiddleware.py @@ -2,12 +2,12 @@ from __future__ import annotations from collections.abc import AsyncIterator, Iterable from inspect import isasyncgen +from typing import Any from unittest import mock import pytest from testfixtures import LogCapture from twisted.internet import defer -from twisted.python.failure import Failure from twisted.trial.unittest import TestCase from scrapy.core.spidermw import SpiderMiddlewareManager @@ -31,53 +31,51 @@ class TestSpiderMiddleware(TestCase): self.spider = self.crawler._create_spider("foo") self.mwman = SpiderMiddlewareManager.from_crawler(self.crawler) - def _scrape_response(self): + async def _scrape_response(self) -> Any: """Execute spider mw manager's scrape_response method and return the result. Raise exception in case of failure. """ scrape_func = mock.MagicMock() - dfd = self.mwman.scrape_response( - scrape_func, self.response, self.request, self.spider + return await maybe_deferred_to_future( + self.mwman.scrape_response( + scrape_func, self.response, self.request, self.spider + ) ) - # catch deferred result and return the value - results = [] - dfd.addBoth(results.append) - self._wait(dfd) - return results[0] class TestProcessSpiderInputInvalidOutput(TestSpiderMiddleware): """Invalid return value for process_spider_input method""" - def test_invalid_process_spider_input(self): + @deferred_f_from_coro_f + async def test_invalid_process_spider_input(self): class InvalidProcessSpiderInputMiddleware: def process_spider_input(self, response, spider): return 1 self.mwman._add_middleware(InvalidProcessSpiderInputMiddleware()) - result = self._scrape_response() - assert isinstance(result, Failure) - assert isinstance(result.value, _InvalidOutput) + with pytest.raises(_InvalidOutput): + await self._scrape_response() class TestProcessSpiderOutputInvalidOutput(TestSpiderMiddleware): """Invalid return value for process_spider_output method""" - def test_invalid_process_spider_output(self): + @deferred_f_from_coro_f + async def test_invalid_process_spider_output(self): class InvalidProcessSpiderOutputMiddleware: def process_spider_output(self, response, result, spider): return 1 self.mwman._add_middleware(InvalidProcessSpiderOutputMiddleware()) - result = self._scrape_response() - assert isinstance(result, Failure) - assert isinstance(result.value, _InvalidOutput) + with pytest.raises(_InvalidOutput): + await self._scrape_response() class TestProcessSpiderExceptionInvalidOutput(TestSpiderMiddleware): """Invalid return value for process_spider_exception method""" - def test_invalid_process_spider_exception(self): + @deferred_f_from_coro_f + async def test_invalid_process_spider_exception(self): class InvalidProcessSpiderOutputExceptionMiddleware: def process_spider_exception(self, response, exception, spider): return 1 @@ -88,15 +86,15 @@ class TestProcessSpiderExceptionInvalidOutput(TestSpiderMiddleware): self.mwman._add_middleware(InvalidProcessSpiderOutputExceptionMiddleware()) self.mwman._add_middleware(RaiseExceptionProcessSpiderOutputMiddleware()) - result = self._scrape_response() - assert isinstance(result, Failure) - assert isinstance(result.value, _InvalidOutput) + with pytest.raises(_InvalidOutput): + await self._scrape_response() class TestProcessSpiderExceptionReRaise(TestSpiderMiddleware): """Re raise the exception by returning None""" - def test_process_spider_exception_return_none(self): + @deferred_f_from_coro_f + async def test_process_spider_exception_return_none(self): class ProcessSpiderExceptionReturnNoneMiddleware: def process_spider_exception(self, response, exception, spider): return None @@ -107,9 +105,8 @@ class TestProcessSpiderExceptionReRaise(TestSpiderMiddleware): self.mwman._add_middleware(ProcessSpiderExceptionReturnNoneMiddleware()) self.mwman._add_middleware(RaiseExceptionProcessSpiderOutputMiddleware()) - result = self._scrape_response() - assert isinstance(result, Failure) - assert isinstance(result.value, ZeroDivisionError) + with pytest.raises(ZeroDivisionError): + await self._scrape_response() class TestBaseAsyncSpiderMiddleware(TestSpiderMiddleware): @@ -350,7 +347,7 @@ class TestProcessStartSimple(TestBaseAsyncSpiderMiddleware): ) self.spider = self.crawler._create_spider() self.mwman = SpiderMiddlewareManager.from_crawler(self.crawler) - return await maybe_deferred_to_future(self.mwman.process_start(self.spider)) + return await self.mwman.process_start(self.spider) @deferred_f_from_coro_f async def test_simple(self): diff --git a/tests/test_spidermiddleware_base.py b/tests/test_spidermiddleware_base.py new file mode 100644 index 000000000..77d055d50 --- /dev/null +++ b/tests/test_spidermiddleware_base.py @@ -0,0 +1,132 @@ +from __future__ import annotations + +from typing import TYPE_CHECKING, Any + +import pytest + +from scrapy import Request, Spider +from scrapy.http import Response +from scrapy.spidermiddlewares.base import BaseSpiderMiddleware +from scrapy.utils.test import get_crawler + +if TYPE_CHECKING: + from scrapy.crawler import Crawler + + +@pytest.fixture +def crawler() -> Crawler: + return get_crawler(Spider) + + +def test_trivial(crawler): + class TrivialSpiderMiddleware(BaseSpiderMiddleware): + pass + + mw = TrivialSpiderMiddleware.from_crawler(crawler) + assert hasattr(mw, "crawler") + assert mw.crawler is crawler + test_req = Request("data:,") + spider_output = [test_req, {"foo": "bar"}] + for processed in [ + list( + mw.process_spider_output(Response("data:,"), spider_output, crawler.spider) + ), + list(mw.process_start_requests(spider_output, crawler.spider)), + ]: + assert processed == [test_req, {"foo": "bar"}] + + +def test_processed_request(crawler): + class ProcessReqSpiderMiddleware(BaseSpiderMiddleware): + def get_processed_request( + self, request: Request, response: Response | None + ) -> Request | None: + if request.url == "data:2,": + return None + if request.url == "data:3,": + return Request("data:30,") + return request + + mw = ProcessReqSpiderMiddleware.from_crawler(crawler) + test_req1 = Request("data:1,") + test_req2 = Request("data:2,") + test_req3 = Request("data:3,") + spider_output = [test_req1, {"foo": "bar"}, test_req2, test_req3] + for processed in [ + list( + mw.process_spider_output(Response("data:,"), spider_output, crawler.spider) + ), + list(mw.process_start_requests(spider_output, crawler.spider)), + ]: + assert len(processed) == 3 + assert isinstance(processed[0], Request) + assert processed[0].url == "data:1," + assert processed[1] == {"foo": "bar"} + assert isinstance(processed[2], Request) + assert processed[2].url == "data:30," + + +def test_processed_item(crawler): + class ProcessItemSpiderMiddleware(BaseSpiderMiddleware): + def get_processed_item(self, item: Any, response: Response | None) -> Any: + if item["foo"] == 2: + return None + if item["foo"] == 3: + item["foo"] = 30 + return item + + mw = ProcessItemSpiderMiddleware.from_crawler(crawler) + test_req = Request("data:,") + spider_output = [{"foo": 1}, {"foo": 2}, test_req, {"foo": 3}] + for processed in [ + list( + mw.process_spider_output(Response("data:,"), spider_output, crawler.spider) + ), + list(mw.process_start_requests(spider_output, crawler.spider)), + ]: + assert processed == [{"foo": 1}, test_req, {"foo": 30}] + + +def test_processed_both(crawler): + class ProcessBothSpiderMiddleware(BaseSpiderMiddleware): + def get_processed_request( + self, request: Request, response: Response | None + ) -> Request | None: + if request.url == "data:2,": + return None + if request.url == "data:3,": + return Request("data:30,") + return request + + def get_processed_item(self, item: Any, response: Response | None) -> Any: + if item["foo"] == 2: + return None + if item["foo"] == 3: + item["foo"] = 30 + return item + + mw = ProcessBothSpiderMiddleware.from_crawler(crawler) + test_req1 = Request("data:1,") + test_req2 = Request("data:2,") + test_req3 = Request("data:3,") + spider_output = [ + test_req1, + {"foo": 1}, + {"foo": 2}, + test_req2, + {"foo": 3}, + test_req3, + ] + for processed in [ + list( + mw.process_spider_output(Response("data:,"), spider_output, crawler.spider) + ), + list(mw.process_start_requests(spider_output, crawler.spider)), + ]: + assert len(processed) == 4 + assert isinstance(processed[0], Request) + assert processed[0].url == "data:1," + assert processed[1] == {"foo": 1} + assert processed[2] == {"foo": 30} + assert isinstance(processed[3], Request) + assert processed[3].url == "data:30," diff --git a/tests/test_spidermiddleware_depth.py b/tests/test_spidermiddleware_depth.py index dfcc141c3..9b4aa624c 100644 --- a/tests/test_spidermiddleware_depth.py +++ b/tests/test_spidermiddleware_depth.py @@ -1,19 +1,18 @@ from scrapy.http import Request, Response from scrapy.spidermiddlewares.depth import DepthMiddleware from scrapy.spiders import Spider -from scrapy.statscollectors import StatsCollector from scrapy.utils.test import get_crawler class TestDepthMiddleware: def setup_method(self): - crawler = get_crawler(Spider) + crawler = get_crawler(Spider, {"DEPTH_LIMIT": 1, "DEPTH_STATS_VERBOSE": True}) self.spider = crawler._create_spider("scrapytest.org") - self.stats = StatsCollector(crawler) + self.stats = crawler.stats self.stats.open_spider(self.spider) - self.mw = DepthMiddleware(1, self.stats, True) + self.mw = DepthMiddleware.from_crawler(crawler) def test_process_spider_output(self): req = Request("http://scrapytest.org") diff --git a/tests/test_spidermiddleware_offsite.py b/tests/test_spidermiddleware_offsite.py index f4563a0a4..e4f4b8f9b 100644 --- a/tests/test_spidermiddleware_offsite.py +++ b/tests/test_spidermiddleware_offsite.py @@ -10,7 +10,7 @@ from scrapy.utils.test import get_crawler class TestOffsiteMiddleware: def setup_method(self): crawler = get_crawler(Spider) - self.spider = crawler._create_spider(**self._get_spiderargs()) + self.spider = crawler.spider = crawler._create_spider(**self._get_spiderargs()) self.mw = OffsiteMiddleware.from_crawler(crawler) self.mw.spider_opened(self.spider) diff --git a/tests/test_spidermiddleware_start.py b/tests/test_spidermiddleware_start.py new file mode 100644 index 000000000..295b10ea8 --- /dev/null +++ b/tests/test_spidermiddleware_start.py @@ -0,0 +1,44 @@ +from twisted.trial.unittest import TestCase + +from scrapy.http import Request +from scrapy.spidermiddlewares.start import StartSpiderMiddleware +from scrapy.spiders import Spider +from scrapy.utils.defer import deferred_f_from_coro_f +from scrapy.utils.misc import build_from_crawler +from scrapy.utils.test import get_crawler + + +class TestMiddleware(TestCase): + @deferred_f_from_coro_f + async def test_async(self): + crawler = get_crawler(Spider) + mw = build_from_crawler(StartSpiderMiddleware, crawler) + + async def start(): + yield Request("data:,1") + yield Request("data:,2", meta={"is_start_request": True}) + yield Request("data:,2", meta={"is_start_request": False}) + yield Request("data:,2", meta={"is_start_request": "foo"}) + + result = [ + request.meta["is_start_request"] + async for request in mw.process_start(start()) + ] + assert result == [True, True, False, "foo"] + + @deferred_f_from_coro_f + async def test_sync(self): + crawler = get_crawler(Spider) + mw = build_from_crawler(StartSpiderMiddleware, crawler) + + def start(): + yield Request("data:,1") + yield Request("data:,2", meta={"is_start_request": True}) + yield Request("data:,2", meta={"is_start_request": False}) + yield Request("data:,2", meta={"is_start_request": "foo"}) + + result = [ + request.meta["is_start_request"] + for request in mw.process_start_requests(start(), Spider("test")) + ] + assert result == [True, True, False, "foo"] diff --git a/tox.ini b/tox.ini index 1406811d9..ba6490b14 100644 --- a/tox.ini +++ b/tox.ini @@ -143,7 +143,7 @@ deps = google-cloud-storage ipython robotexclusionrulesparser - uvloop; platform_system != "Windows" + uvloop; platform_system != "Windows" and implementation_name != "pypy" zstandard; implementation_name != "pypy" # optional for HTTP compress downloader middleware tests [testenv:extra-deps-pinned] @@ -159,7 +159,7 @@ deps = google-cloud-storage==1.29.0 ipython==2.0.0 robotexclusionrulesparser==1.6.2 - uvloop==0.14.0; platform_system != "Windows" + uvloop==0.14.0; platform_system != "Windows" and implementation_name != "pypy" zstandard==0.1; implementation_name != "pypy" install_command = {[pinned]install_command} setenv =