diff --git a/scrapy/commands/parse.py b/scrapy/commands/parse.py index 0dd9954cb..c4b3d2af9 100644 --- a/scrapy/commands/parse.py +++ b/scrapy/commands/parse.py @@ -282,6 +282,7 @@ class Command(BaseRunSpiderCommand): ) -> list[Any]: items, requests, opts, depth, spider, callback = args if opts.pipelines: + assert self.pcrawler.engine itemproc = self.pcrawler.engine.scraper.itemproc for item in items: itemproc.process_item(item, spider) diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index 78dc16df6..5468398aa 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -5,29 +5,32 @@ import warnings from collections import deque from datetime import datetime from time import time -from typing import TYPE_CHECKING, Any, TypeVar, cast +from typing import TYPE_CHECKING, Any, cast from twisted.internet import task -from twisted.internet.defer import Deferred +from twisted.internet.defer import Deferred, inlineCallbacks from scrapy import Request, Spider, signals from scrapy.core.downloader.handlers import DownloadHandlers from scrapy.core.downloader.middleware import DownloaderMiddlewareManager from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.resolver import dnscache -from scrapy.utils.defer import mustbe_deferred +from scrapy.utils.defer import ( + deferred_from_coro, + maybe_deferred_to_future, + mustbe_deferred, +) from scrapy.utils.httpobj import urlparse_cached if TYPE_CHECKING: + from collections.abc import Generator + from scrapy.crawler import Crawler from scrapy.http import Response from scrapy.settings import BaseSettings from scrapy.signalmanager import SignalManager -_T = TypeVar("_T") - - class Slot: """Downloader slot""" @@ -114,16 +117,17 @@ class Downloader: "DOWNLOAD_SLOTS", {} ) - def fetch(self, request: Request, spider: Spider) -> Deferred[Response | Request]: - def _deactivate(response: _T) -> _T: - self.active.remove(request) - return response - + @inlineCallbacks + def fetch( + self, request: Request, spider: Spider + ) -> Generator[Deferred[Any], Any, Response | Request]: self.active.add(request) - dfd: Deferred[Response | Request] = self.middleware.download( - self._enqueue_request, request, spider - ) - return dfd.addBoth(_deactivate) + try: + return ( + yield self.middleware.download(self._enqueue_request, request, spider) + ) + finally: + self.active.remove(request) def needs_backout(self) -> bool: return len(self.active) >= self.total_concurrency @@ -164,22 +168,23 @@ class Downloader: ) return self.get_slot_key(request) - def _enqueue_request(self, request: Request, spider: Spider) -> Deferred[Response]: + @inlineCallbacks + def _enqueue_request( + self, request: Request, spider: Spider + ) -> Generator[Deferred[Any], Any, Response]: key, slot = self._get_slot(request, spider) request.meta[self.DOWNLOAD_SLOT] = key - - def _deactivate(response: Response) -> Response: - slot.active.remove(request) - return response - slot.active.add(request) self.signals.send_catch_log( signal=signals.request_reached_downloader, request=request, spider=spider ) - deferred: Deferred[Response] = Deferred().addBoth(_deactivate) - slot.queue.append((request, deferred)) + d: Deferred[Response] = Deferred() + slot.queue.append((request, d)) self._process_queue(spider, slot) - return deferred + try: + return (yield d) + finally: + slot.active.remove(request) def _process_queue(self, spider: Spider, slot: Slot) -> None: from twisted.internet import reactor @@ -202,26 +207,23 @@ class Downloader: while slot.queue and slot.free_transfer_slots() > 0: slot.lastseen = now request, deferred = slot.queue.popleft() - dfd = self._download(slot, request, spider) + dfd = deferred_from_coro(self._download(slot, request, spider)) dfd.chainDeferred(deferred) # prevent burst if inter-request delays were configured if delay: self._process_queue(spider, slot) break - def _download( - self, slot: Slot, request: Request, spider: Spider - ) -> Deferred[Response]: - # The order is very important for the following deferreds. Do not change! - - # 1. Create the download deferred - dfd: Deferred[Response] = mustbe_deferred( - self.handlers.download_request, request, spider - ) - - # 2. Notify response_downloaded listeners about the recent download - # before querying queue for next request - def _downloaded(response: Response) -> Response: + async def _download(self, slot: Slot, request: Request, spider: Spider) -> Response: + # The order is very important for the following logic. Do not change! + slot.transferring.add(request) + try: + # 1. Download the response + response: Response = await maybe_deferred_to_future( + mustbe_deferred(self.handlers.download_request, request, spider) + ) + # 2. Notify response_downloaded listeners about the recent download + # before querying queue for next request self.signals.send_catch_log( signal=signals.response_downloaded, response=response, @@ -229,24 +231,16 @@ class Downloader: spider=spider, ) return response - - dfd.addCallback(_downloaded) - - # 3. After response arrives, remove the request from transferring - # state to free up the transferring slot so it can be used by the - # following requests (perhaps those which came from the downloader - # middleware itself) - slot.transferring.add(request) - - def finish_transferring(_: _T) -> _T: + finally: + # 3. After response arrives, remove the request from transferring + # state to free up the transferring slot so it can be used by the + # following requests (perhaps those which came from the downloader + # middleware itself) slot.transferring.remove(request) self._process_queue(spider, slot) self.signals.send_catch_log( signal=signals.request_left_downloader, request=request, spider=spider ) - return _ - - return dfd.addBoth(finish_transferring) def close(self) -> None: self._slot_gc_loop.stop() diff --git a/scrapy/core/downloader/middleware.py b/scrapy/core/downloader/middleware.py index db4191385..a4055849d 100644 --- a/scrapy/core/downloader/middleware.py +++ b/scrapy/core/downloader/middleware.py @@ -20,8 +20,6 @@ from scrapy.utils.defer import deferred_from_coro, mustbe_deferred if TYPE_CHECKING: from collections.abc import Generator - from twisted.python.failure import Failure - from scrapy import Spider from scrapy.settings import BaseSettings @@ -41,12 +39,13 @@ class DownloaderMiddlewareManager(MiddlewareManager): if hasattr(mw, "process_exception"): self.methods["process_exception"].appendleft(mw.process_exception) + @inlineCallbacks def download( self, download_func: Callable[[Request, Spider], Deferred[Response]], request: Request, spider: Spider, - ) -> Deferred[Response | Request]: + ) -> Generator[Deferred[Any], Any, Response | Request]: @inlineCallbacks def process_request( request: Request, @@ -92,9 +91,8 @@ class DownloaderMiddlewareManager(MiddlewareManager): @inlineCallbacks def process_exception( - failure: Failure, - ) -> Generator[Deferred[Any], Any, Failure | Response | Request]: - exception = failure.value + exception: Exception, + ) -> Generator[Deferred[Any], Any, Response | Request]: for method in self.methods["process_exception"]: method = cast(Callable, method) response = yield deferred_from_coro( @@ -109,11 +107,12 @@ class DownloaderMiddlewareManager(MiddlewareManager): ) if response: return response - return failure + raise exception - deferred: Deferred[Response | Request] = mustbe_deferred( - process_request, request - ) - deferred.addErrback(process_exception) - deferred.addCallback(process_response) - return deferred + try: + result: Response | Request = yield mustbe_deferred(process_request, request) + except Exception as ex: + # either returns a request or response (which we pass to process_response()) + # or reraises the exception + result = yield process_exception(ex) + return (yield process_response(result)) diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 7f5dd0405..658f6e774 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -10,7 +10,7 @@ from __future__ import annotations import logging from time import time from traceback import format_exc -from typing import TYPE_CHECKING, Any, TypeVar, cast +from typing import TYPE_CHECKING, Any, cast from twisted.internet.defer import Deferred, inlineCallbacks, succeed from twisted.internet.task import LoopingCall @@ -42,8 +42,6 @@ if TYPE_CHECKING: logger = logging.getLogger(__name__) -_T = TypeVar("_T") - class _Slot: def __init__( @@ -349,28 +347,32 @@ class ExecutionEngine: signals.request_dropped, request=request, spider=self.spider ) - def download(self, request: Request) -> Deferred[Response]: + @inlineCallbacks + def download(self, request: Request) -> Generator[Deferred[Any], Any, Response]: """Return a Deferred which fires with a Response as result, only downloader middlewares are applied""" if self.spider is None: raise RuntimeError(f"No open spider to crawl: {request}") - d: Deferred[Response | Request] = self._download(request) - # Deferred.addBoth() overloads don't seem to support a Union[_T, Deferred[_T]] return type - d2: Deferred[Response] = d.addBoth(self._downloaded, request) # type: ignore[call-overload] - return d2 + try: + response_or_request = yield self._download(request) + finally: + assert self._slot is not None + self._slot.remove_request(request) + if isinstance(response_or_request, Request): + return (yield self.download(response_or_request)) + return response_or_request - def _downloaded( - self, result: Response | Request | Failure, request: Request - ) -> Deferred[Response] | Response | Failure: - assert self._slot is not None # typing - self._slot.remove_request(request) - return self.download(result) if isinstance(result, Request) else result - - def _download(self, request: Request) -> Deferred[Response | Request]: + @inlineCallbacks + def _download( + self, request: Request + ) -> Generator[Deferred[Any], Any, Response | Request]: assert self._slot is not None # typing + assert self.spider is not None self._slot.add_request(request) - - def _on_success(result: Response | Request) -> Response | Request: + try: + result: Response | Request = yield self.downloader.fetch( + request, self.spider + ) if not isinstance(result, (Response, Request)): raise TypeError( f"Incorrect type: expected Response or Request, got {type(result)}: {result!r}" @@ -391,17 +393,8 @@ class ExecutionEngine: spider=self.spider, ) return result - - def _on_complete(_: _T) -> _T: - assert self._slot is not None + finally: self._slot.nextcall.schedule() - return _ - - assert self.spider is not None - dwld: Deferred[Response | Request] = self.downloader.fetch(request, self.spider) - dwld.addCallback(_on_success) - dwld.addBoth(_on_complete) - return dwld @deferred_f_from_coro_f async def open_spider( diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index 9378f2651..2c48a9a81 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -111,11 +111,11 @@ class Scraper: assert crawler.logformatter self.logformatter: LogFormatter = crawler.logformatter - @inlineCallbacks - def open_spider(self, spider: Spider) -> Generator[Deferred[Any], Any, None]: + @deferred_f_from_coro_f + async def open_spider(self, spider: Spider) -> None: """Open the given spider for scraping and allocate resources for it""" self.slot = Slot(self.crawler.settings.getint("SCRAPER_SLOT_MAX_ACTIVE_SIZE")) - yield self.itemproc.open_spider(spider) + await maybe_deferred_to_future(self.itemproc.open_spider(spider)) def close_spider(self, spider: Spider | None = None) -> Deferred[Spider]: """Close a spider being scraped and release its resources""" @@ -191,10 +191,8 @@ class Scraper: if isinstance(result, Response): 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 - ) + output = await self.spidermw.scrape_response_async( + self.call_spider, result, request, self.crawler.spider ) except Exception: self.handle_spider_error(Failure(), request, result) @@ -363,12 +361,19 @@ class Scraper: self.crawler.engine.crawl(request=output) return if output is not None: - await maybe_deferred_to_future( - self.start_itemproc(output, response=response) - ) + await self.start_itemproc_async(output, response=response) - @deferred_f_from_coro_f - async def start_itemproc(self, item: Any, *, response: Response | None) -> None: + def start_itemproc(self, item: Any, *, response: Response | None) -> Deferred[None]: + """Send *item* to the item pipelines for processing. + + *response* is the source of the item data. If the item does not come + from response data, e.g. it was hard-coded, set it to ``None``. + """ + return deferred_from_coro(self.start_itemproc_async(item, response=response)) + + async def start_itemproc_async( + 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 diff --git a/scrapy/core/spidermw.py b/scrapy/core/spidermw.py index 310abb9b7..10aad7858 100644 --- a/scrapy/core/spidermw.py +++ b/scrapy/core/spidermw.py @@ -23,7 +23,6 @@ from scrapy.middleware import MiddlewareManager from scrapy.utils.asyncgen import as_async_generator, collect_asyncgen from scrapy.utils.conf import build_component_list from scrapy.utils.defer import ( - deferred_f_from_coro_f, deferred_from_coro, maybe_deferred_to_future, mustbe_deferred, @@ -169,7 +168,7 @@ class SpiderMiddlewareManager(MiddlewareManager): exception_result = cast( Union[Failure, MutableChain[_T]], self._process_spider_exception( - response, spider, Failure(ex), exception_processor_index + response, spider, ex, exception_processor_index ), ) if isinstance(exception_result, Failure): @@ -185,7 +184,7 @@ class SpiderMiddlewareManager(MiddlewareManager): exception_result = cast( Union[Failure, MutableAsyncChain[_T]], self._process_spider_exception( - response, spider, Failure(ex), exception_processor_index + response, spider, ex, exception_processor_index ), ) if isinstance(exception_result, Failure): @@ -201,13 +200,12 @@ class SpiderMiddlewareManager(MiddlewareManager): self, response: Response, spider: Spider, - _failure: Failure, + exception: Exception, start_index: int = 0, - ) -> Failure | MutableChain[_T] | MutableAsyncChain[_T]: - exception = _failure.value + ) -> MutableChain[_T] | MutableAsyncChain[_T]: # don't handle _InvalidOutput exception if isinstance(exception, _InvalidOutput): - return _failure + raise exception method_list = islice( self.methods["process_spider_exception"], start_index, None ) @@ -242,7 +240,7 @@ class SpiderMiddlewareManager(MiddlewareManager): f"or an iterable, got {type(result)}" ) raise _InvalidOutput(msg) - return _failure + raise exception # This method cannot be made async def, as _process_spider_exception relies on the Deferred result # being available immediately which doesn't work when it's a wrapped coroutine. @@ -308,7 +306,7 @@ class SpiderMiddlewareManager(MiddlewareManager): except Exception as ex: exception_result: Failure | MutableChain[_T] | MutableAsyncChain[_T] = ( self._process_spider_exception( - response, spider, Failure(ex), method_index + 1 + response, spider, ex, method_index + 1 ) ) if isinstance(exception_result, Failure): @@ -369,24 +367,36 @@ class SpiderMiddlewareManager(MiddlewareManager): request: Request, spider: Spider, ) -> Deferred[MutableChain[_T] | MutableAsyncChain[_T]]: + return deferred_from_coro( + self.scrape_response_async(scrape_func, response, request, spider) + ) + + async def scrape_response_async( + self, + scrape_func: ScrapeFunc[_T], + response: Response, + request: Request, + spider: Spider, + ) -> MutableChain[_T] | MutableAsyncChain[_T]: async def process_callback_output( result: Iterable[_T] | AsyncIterator[_T], ) -> MutableChain[_T] | MutableAsyncChain[_T]: return await self._process_callback_output(response, spider, result) def process_spider_exception( - _failure: Failure, - ) -> Failure | MutableChain[_T] | MutableAsyncChain[_T]: - return self._process_spider_exception(response, spider, _failure) + exception: Exception, + ) -> MutableChain[_T] | MutableAsyncChain[_T]: + return self._process_spider_exception(response, spider, exception) - dfd: Deferred[Iterable[_T] | AsyncIterator[_T]] = mustbe_deferred( - self._process_spider_input, scrape_func, response, request, spider - ) - dfd2: Deferred[MutableChain[_T] | MutableAsyncChain[_T]] = dfd.addCallback( - deferred_f_from_coro_f(process_callback_output) - ) - dfd2.addErrback(process_spider_exception) - return dfd2 + try: + it: Iterable[_T] | AsyncIterator[_T] = await maybe_deferred_to_future( + mustbe_deferred( + self._process_spider_input, scrape_func, response, request, spider + ) + ) + return await process_callback_output(it) + except Exception as ex: + return process_spider_exception(ex) async def process_start(self, spider: Spider) -> AsyncIterator[Any] | None: self._check_deprecated_start_requests_use(spider) diff --git a/scrapy/crawler.py b/scrapy/crawler.py index 749096db5..5dbee6537 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -10,7 +10,6 @@ from twisted.internet.defer import ( Deferred, DeferredList, inlineCallbacks, - maybeDeferred, ) from zope.interface.verify import verifyClass @@ -175,7 +174,7 @@ class Crawler: if self.crawling: self.crawling = False assert self.engine - yield maybeDeferred(self.engine.stop) + yield self.engine.stop() @staticmethod def _get_component( @@ -277,12 +276,6 @@ class CrawlerRunner: process. See :ref:`run-from-script` for an example. """ - crawlers = property( - lambda self: self._crawlers, - doc="Set of :class:`crawlers ` started by " - ":meth:`crawl` and managed by this class.", - ) - @staticmethod def _get_spider_loader(settings: BaseSettings) -> SpiderLoaderProtocol: """Get SpiderLoader instance from settings""" @@ -303,6 +296,12 @@ class CrawlerRunner: self._active: set[Deferred[None]] = set() self.bootstrap_failed = False + @property + def crawlers(self) -> set[Crawler]: + """Set of :class:`crawlers ` started by + :meth:`crawl` and managed by this class.""" + return self._crawlers + def crawl( self, crawler_or_spidercls: type[Spider] | str | Crawler, @@ -338,18 +337,19 @@ class CrawlerRunner: crawler = self.create_crawler(crawler_or_spidercls) return self._crawl(crawler, *args, **kwargs) - def _crawl(self, crawler: Crawler, *args: Any, **kwargs: Any) -> Deferred[None]: + @inlineCallbacks + def _crawl( + self, crawler: Crawler, *args: Any, **kwargs: Any + ) -> Generator[Deferred[Any], Any, None]: self.crawlers.add(crawler) d = crawler.crawl(*args, **kwargs) self._active.add(d) - - def _done(result: _T) -> _T: + try: + yield d + finally: self.crawlers.discard(crawler) self._active.discard(d) self.bootstrap_failed |= not getattr(crawler, "spider", None) - return result - - return d.addBoth(_done) def create_crawler( self, crawler_or_spidercls: type[Spider] | str | Crawler @@ -501,10 +501,12 @@ class CrawlerProcess(CrawlerRunner): ) reactor.run(installSignalHandlers=install_signal_handlers) # blocking call - def _graceful_stop_reactor(self) -> Deferred[Any]: - d = self.stop() - d.addBoth(self._stop_reactor) - return d + @inlineCallbacks + def _graceful_stop_reactor(self) -> Generator[Deferred[Any], Any, None]: + try: + yield self.stop() + finally: + self._stop_reactor() def _stop_reactor(self, _: Any = None) -> None: from twisted.internet import reactor diff --git a/tests/test_downloadermiddleware.py b/tests/test_downloadermiddleware.py index 8ae160f8a..61a5a7df5 100644 --- a/tests/test_downloadermiddleware.py +++ b/tests/test_downloadermiddleware.py @@ -131,6 +131,39 @@ class TestResponseFromProcessRequest(TestManagerBase): assert not download_func.called +class TestResponseFromProcessException(TestManagerBase): + """Tests middleware returning a response from process_exception.""" + + @deferred_f_from_coro_f + async def test_process_response_called(self): + resp = Response("http://example.com/index.html") + calls = [] + + def download_func(request, spider): + raise ValueError("test") + + class ResponseMiddleware: + def process_response(self, request, response, spider): + calls.append("process_response") + return resp + + def process_exception(self, request, exception, spider): + calls.append("process_exception") + return resp + + self.mwman._add_middleware(ResponseMiddleware()) + + req = Request("http://example.com/index.html") + result = await maybe_deferred_to_future( + self.mwman.download(download_func, req, self.spider) + ) + assert result is resp + assert calls == [ + "process_exception", + "process_response", + ] + + class TestInvalidOutput(TestManagerBase): @deferred_f_from_coro_f async def test_invalid_process_request(self):