Refactor more Deferred functions. (#6795)

This commit is contained in:
Andrey Rakhmatullin 2025-05-13 22:47:57 +04:00 committed by GitHub
parent 2442536d0f
commit b86f00327a
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
8 changed files with 179 additions and 142 deletions

View File

@ -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)

View File

@ -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()

View File

@ -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))

View File

@ -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(

View File

@ -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

View File

@ -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)

View File

@ -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 <scrapy.crawler.Crawler>` 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 <scrapy.crawler.Crawler>` 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

View File

@ -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):