diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index 501c669ce..9293d7b78 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -14,7 +14,12 @@ 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.asyncio import AsyncioLoopingCall, create_looping_call +from scrapy.utils.asyncio import ( + AsyncioLoopingCall, + CallLaterResult, + call_later, + create_looping_call, +) from scrapy.utils.defer import ( deferred_from_coro, maybe_deferred_to_future, @@ -50,7 +55,7 @@ class Slot: self.queue: deque[tuple[Request, Deferred[Response]]] = deque() self.transferring: set[Request] = set() self.lastseen: float = 0 - self.latercall = None + self.latercall: CallLaterResult | None = None def free_transfer_slots(self) -> int: return self.concurrency - len(self.transferring) @@ -61,8 +66,9 @@ class Slot: return self.delay def close(self) -> None: - if self.latercall and self.latercall.active(): + if self.latercall: self.latercall.cancel() + self.latercall = None def __repr__(self) -> str: cls_name = self.__class__.__name__ @@ -191,9 +197,8 @@ class Downloader: slot.active.remove(request) def _process_queue(self, spider: Spider, slot: Slot) -> None: - from twisted.internet import reactor - - if slot.latercall and slot.latercall.active(): + if slot.latercall: + # block processing until slot.latercall is called return # Delay queue processing if a download_delay is configured @@ -202,9 +207,7 @@ class Downloader: if delay: penalty = delay - now + slot.lastseen if penalty > 0: - slot.latercall = reactor.callLater( - penalty, self._process_queue, spider, slot - ) + slot.latercall = call_later(penalty, self._latercall, spider, slot) return # Process enqueued requests if there are free slots to transfer for this slot @@ -218,6 +221,10 @@ class Downloader: self._process_queue(spider, slot) break + def _latercall(self, spider: Spider, slot: Slot) -> None: + slot.latercall = None + self._process_queue(spider, slot) + 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) diff --git a/scrapy/extensions/closespider.py b/scrapy/extensions/closespider.py index a649a86e2..b4c6c73a0 100644 --- a/scrapy/extensions/closespider.py +++ b/scrapy/extensions/closespider.py @@ -12,9 +12,15 @@ from typing import TYPE_CHECKING, Any from scrapy import Request, Spider, signals from scrapy.exceptions import NotConfigured -from scrapy.utils.asyncio import create_looping_call +from scrapy.utils.asyncio import ( + AsyncioLoopingCall, + CallLaterResult, + call_later, + create_looping_call, +) if TYPE_CHECKING: + from twisted.internet.task import LoopingCall from twisted.python.failure import Failure # typing.Self requires Python 3.11 @@ -31,6 +37,12 @@ class CloseSpider: def __init__(self, crawler: Crawler): self.crawler: Crawler = crawler + # for CLOSESPIDER_TIMEOUT + self.task: CallLaterResult | None = None + + # for CLOSESPIDER_TIMEOUT_NO_ITEM + self.task_no_item: AsyncioLoopingCall | LoopingCall | None = None + self.close_on: dict[str, Any] = { "timeout": crawler.settings.getfloat("CLOSESPIDER_TIMEOUT"), "itemcount": crawler.settings.getint("CLOSESPIDER_ITEMCOUNT"), @@ -92,14 +104,12 @@ class CloseSpider: self.crawler.engine.close_spider(spider, "closespider_pagecount_no_item") def spider_opened(self, spider: Spider) -> None: - from twisted.internet import reactor - assert self.crawler.engine - self.task = reactor.callLater( + self.task = call_later( self.close_on["timeout"], self.crawler.engine.close_spider, spider, - reason="closespider_timeout", + "closespider_timeout", ) def item_scraped(self, item: Any, spider: Spider) -> None: @@ -110,13 +120,14 @@ class CloseSpider: self.crawler.engine.close_spider(spider, "closespider_itemcount") def spider_closed(self, spider: Spider) -> None: - task = getattr(self, "task", None) - if task and task.active(): - task.cancel() + if self.task: + self.task.cancel() + self.task = None - task_no_item = getattr(self, "task_no_item", None) - if task_no_item and task_no_item.running: - task_no_item.stop() + if self.task_no_item: + if self.task_no_item.running: + self.task_no_item.stop() + self.task_no_item = None def spider_opened_no_item(self, spider: Spider) -> None: self.task_no_item = create_looping_call(self._count_items_produced, spider) diff --git a/scrapy/utils/asyncio.py b/scrapy/utils/asyncio.py index cae2dc033..8c5b843cb 100644 --- a/scrapy/utils/asyncio.py +++ b/scrapy/utils/asyncio.py @@ -15,10 +15,14 @@ from scrapy.utils.asyncgen import as_async_generator from scrapy.utils.reactor import is_asyncio_reactor_installed, is_reactor_installed if TYPE_CHECKING: + from twisted.internet.base import DelayedCall + # typing.Concatenate and typing.ParamSpec require Python 3.10 - from typing_extensions import Concatenate, ParamSpec + # typing.Self, typing.TypeVarTuple and typing.Unpack require Python 3.11 + from typing_extensions import Concatenate, ParamSpec, Self, TypeVarTuple, Unpack _P = ParamSpec("_P") + _Ts = TypeVarTuple("_Ts") _T = TypeVar("_T") @@ -192,3 +196,60 @@ def create_looping_call( if is_asyncio_available(): return AsyncioLoopingCall(func, *args, **kwargs) return LoopingCall(func, *args, **kwargs) + + +def call_later( + delay: float, func: Callable[[Unpack[_Ts]], object], *args: Unpack[_Ts] +) -> CallLaterResult: + """Schedule a function to be called after a delay. + + This uses either ``loop.call_later()`` or ``reactor.callLater()``, depending + on whether asyncio support is available. + """ + if is_asyncio_available(): + loop = asyncio.get_event_loop() + return CallLaterResult.from_asyncio(loop.call_later(delay, func, *args)) + + from twisted.internet import reactor + + return CallLaterResult.from_twisted(reactor.callLater(delay, func, *args)) + + +class CallLaterResult: + """An universal result for :func:`call_later`, wrapping either + :class:`asyncio.TimerHandle` or :class:`twisted.internet.base.DelayedCall`. + + The provided API is close to the :class:`asyncio.TimerHandle` one: there is + no ``active()`` (as there is no such public API in + :class:`asyncio.TimerHandle`) but ``cancel()`` can be called on already + called or cancelled instances. + """ + + _timer_handle: asyncio.TimerHandle | None = None + _delayed_call: DelayedCall | None = None + + @classmethod + def from_asyncio(cls, timer_handle: asyncio.TimerHandle) -> Self: + """Create a CallLaterResult from an asyncio TimerHandle.""" + o = cls() + o._timer_handle = timer_handle + return o + + @classmethod + def from_twisted(cls, delayed_call: DelayedCall) -> Self: + """Create a CallLaterResult from a Twisted DelayedCall.""" + o = cls() + o._delayed_call = delayed_call + return o + + def cancel(self) -> None: + """Cancel the underlying delayed call. + + Does nothing if the delayed call was already called or cancelled. + """ + if self._timer_handle: + self._timer_handle.cancel() + self._timer_handle = None + elif self._delayed_call and self._delayed_call.active(): + self._delayed_call.cancel() + self._delayed_call = None diff --git a/scrapy/utils/reactor.py b/scrapy/utils/reactor.py index 2fb1e0ce7..76f42392b 100644 --- a/scrapy/utils/reactor.py +++ b/scrapy/utils/reactor.py @@ -16,13 +16,14 @@ if TYPE_CHECKING: from asyncio import AbstractEventLoop, AbstractEventLoopPolicy from collections.abc import Callable - from twisted.internet.base import DelayedCall from twisted.internet.protocol import ServerFactory from twisted.internet.tcp import Port # typing.ParamSpec requires Python 3.10 from typing_extensions import ParamSpec + from scrapy.utils.asyncio import CallLaterResult + _P = ParamSpec("_P") _T = TypeVar("_T") @@ -55,27 +56,27 @@ class CallLaterOnce(Generic[_T]): self._func: Callable[_P, _T] = func self._a: tuple[Any, ...] = a self._kw: dict[str, Any] = kw - self._call: DelayedCall | None = None + self._call: CallLaterResult | None = None self._deferreds: list[Deferred] = [] def schedule(self, delay: float = 0) -> None: - from twisted.internet import reactor + from scrapy.utils.asyncio import call_later if self._call is None: - self._call = reactor.callLater(delay, self) + self._call = call_later(delay, self) def cancel(self) -> None: if self._call: self._call.cancel() def __call__(self) -> _T: - from twisted.internet import reactor + from scrapy.utils.asyncio import call_later self._call = None result = self._func(*self._a, **self._kw) for d in self._deferreds: - reactor.callLater(0, d.callback, None) + call_later(0, d.callback, None) self._deferreds = [] return result