diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index d368d7f92..3c7960b86 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -8,6 +8,7 @@ from time import time from typing import TYPE_CHECKING, Any, cast from twisted.internet.defer import Deferred, inlineCallbacks +from twisted.python.failure import Failure from scrapy import Request, Spider, signals from scrapy.core.downloader.handlers import DownloadHandlers @@ -22,7 +23,7 @@ from scrapy.utils.asyncio import ( ) from scrapy.utils.defer import ( _defer_sleep_async, - deferred_from_coro, + _schedule_coro, maybe_deferred_to_future, ) from scrapy.utils.httpobj import urlparse_cached @@ -200,6 +201,7 @@ class Downloader: ) return self.get_slot_key(request) + # passed as download_func into self.middleware.download() in self.fetch() @inlineCallbacks def _enqueue_request( self, request: Request @@ -216,7 +218,7 @@ class Downloader: slot.queue.append((request, d)) self._process_queue(slot) try: - return (yield d) + return (yield d) # fired in _wait_for_download() finally: slot.active.remove(request) @@ -237,9 +239,8 @@ class Downloader: # Process enqueued requests if there are free slots to transfer for this slot while slot.queue and slot.free_transfer_slots() > 0: slot.lastseen = now - request, deferred = slot.queue.popleft() - dfd = deferred_from_coro(self._download(slot, request)) - dfd.chainDeferred(deferred) + request, queue_dfd = slot.queue.popleft() + _schedule_coro(self._wait_for_download(slot, request, queue_dfd)) # prevent burst if inter-request delays were configured if delay: self._process_queue(slot) @@ -282,6 +283,16 @@ class Downloader: spider=self.crawler.spider, ) + async def _wait_for_download( + self, slot: Slot, request: Request, queue_dfd: Deferred[Response] + ) -> None: + try: + response = await self._download(slot, request) + except Exception: + queue_dfd.errback(Failure()) + else: + queue_dfd.callback(response) # awaited in _enqueue_request() + def close(self) -> None: self._slot_gc_loop.stop() for slot in self.slots.values(): diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index 3576c4b5c..c5095f71a 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -27,7 +27,6 @@ from scrapy.utils.defer import ( _defer_sleep_async, _schedule_coro, aiter_errback, - deferred_f_from_coro_f, deferred_from_coro, ensure_awaitable, iter_errback, @@ -235,7 +234,7 @@ class Scraper: dfd = self.slot.add_response_request(result, request) self._scrape_next() try: - yield dfd + yield dfd # fired in _wait_for_processing() except Exception: logger.error( "Scraper bug processing %(request)s", @@ -251,10 +250,9 @@ class Scraper: def _scrape_next(self) -> None: assert self.slot is not None # typing while self.slot.queue: - result, request, deferred = self.slot.next_response_request_deferred() - self._scrape(result, request).chainDeferred(deferred) + result, request, queue_dfd = self.slot.next_response_request_deferred() + _schedule_coro(self._wait_for_processing(result, request, queue_dfd)) - @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)): @@ -296,6 +294,16 @@ class Scraper: else: await self.handle_spider_output_async(output, request, result) + async def _wait_for_processing( + self, result: Response | Failure, request: Request, queue_dfd: Deferred[None] + ) -> None: + try: + await self._scrape(result, request) + except Exception: + queue_dfd.errback(Failure()) + else: + queue_dfd.callback(None) # awaited in enqueue_scrape() + def call_spider( self, result: Response | Failure, request: Request, spider: Spider | None = None ) -> Deferred[Iterable[Any] | AsyncIterator[Any]]: diff --git a/tests/test_core_scraper.py b/tests/test_core_scraper.py new file mode 100644 index 000000000..c819e246e --- /dev/null +++ b/tests/test_core_scraper.py @@ -0,0 +1,27 @@ +from __future__ import annotations + +from typing import TYPE_CHECKING + +from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future +from scrapy.utils.test import get_crawler +from tests.spiders import SimpleSpider + +if TYPE_CHECKING: + import pytest + + from tests.mockserver.http import MockServer + + +@deferred_f_from_coro_f +async def test_scraper_exception( + mockserver: MockServer, + caplog: pytest.LogCaptureFixture, + monkeypatch: pytest.MonkeyPatch, +) -> None: + crawler = get_crawler(SimpleSpider) + monkeypatch.setattr( + "scrapy.core.engine.Scraper.handle_spider_output_async", + lambda *args, **kwargs: 1 / 0, + ) + await maybe_deferred_to_future(crawler.crawl(url=mockserver.url("/"))) + assert "Scraper bug processing" in caplog.text