From a5dc229508e653cc45ff707c6aa2b6a261e6aac9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Sat, 15 Mar 2025 20:50:21 +0100 Subject: [PATCH] =?UTF-8?q?Improve=20exception=20handling=20for=20Schedule?= =?UTF-8?q?r=E2=80=99s=20has=5Fpending=5Frequests()=20and=20enqueue=5Frequ?= =?UTF-8?q?est()?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scrapy/core/engine.py | 24 ++++++++++++-- tests/test_engine_loop.py | 67 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 88 insertions(+), 3 deletions(-) diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 6eaf6e9f6..21b9fe939 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -249,6 +249,16 @@ class ExecutionEngine: and not self._needs_backout() ) + def _scheduler_has_pending_requests(self) -> bool: + try: + return self._slot.scheduler.has_pending_requests() + except Exception as exception: + exception_traceback = format_exc() + logger.error( + f"{global_object_name(self._slot.scheduler.has_pending_requests)} raised an exception: {exception}.\n{exception_traceback}" + ) + return False + @inlineCallbacks def _run_loop(self) -> Generator[Deferred[Any], Any, None]: """Sends new requests from seeds or from the scheduler based on the @@ -288,7 +298,7 @@ class ExecutionEngine: self._spider_idle() elif self._needs_backout(): self._back_in_seconds = self._MIN_BACK_IN_SECONDS - elif self._slot.scheduler.has_pending_requests(): + elif self._scheduler_has_pending_requests(): # If the scheduler reports having pending requests but did not # actually return one, use exponential backoff to schedule a new # call to this method, to see if the scheduler finally returns a @@ -391,7 +401,7 @@ class ExecutionEngine: return False if self._seeds is not None: # not all start requests are handled return False - return not self._slot.scheduler.has_pending_requests() + return not self._scheduler_has_pending_requests() def crawl(self, request: Request) -> None: """Inject the request into the spider <-> downloader pipeline""" @@ -410,7 +420,15 @@ class ExecutionEngine: for handler, result in request_scheduled_result: if isinstance(result, Failure) and isinstance(result.value, IgnoreRequest): return - if not self._slot.scheduler.enqueue_request(request): # type: ignore[union-attr] + try: + request_was_enqueued = self._slot.scheduler.enqueue_request(request) + except Exception as exception: + exception_traceback = format_exc() + logger.error( + f"{global_object_name(self._slot.scheduler.enqueue_request)} raised an exception: {exception}\n{exception_traceback}" + ) + request_was_enqueued = False + if not request_was_enqueued: self.signals.send_catch_log( signals.request_dropped, request=request, spider=spider ) diff --git a/tests/test_engine_loop.py b/tests/test_engine_loop.py index 7fb6298be..15792bc2b 100644 --- a/tests/test_engine_loop.py +++ b/tests/test_engine_loop.py @@ -309,6 +309,73 @@ class MainTestCase(TestCase): expected_urls = ["data:,a", "data:,b", "data:,c"] assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" + # Unexpected scheduler exceptions + + @deferred_f_from_coro_f + async def test_scheduler_has_pending_requests_exception(self): + """If Scheduler.has_pending_requests() raises an exception while + checking if the spider is idle, consider the return value to be False + (i.e. the spider is indeed idle), and log a traceback.""" + + class TestScheduler(MemoryScheduler): + def has_pending_requests(self): + raise RuntimeError + + def next_request(self): + return None + + class TestSpider(Spider): + name = "test" + start_urls = [] + + def parse(self, response): + pass + + actual_urls = [] + + def track_url(request, spider): + actual_urls.append(request.url) + + settings = {"SCHEDULER": TestScheduler} + crawler = get_crawler(TestSpider, settings_dict=settings) + crawler.signals.connect(track_url, signals.request_reached_downloader) + with LogCapture() as log: + await maybe_deferred_to_future(crawler.crawl()) + assert crawler.stats.get_value("finish_reason") == "finished" + expected_urls = [] + assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" + assert "in has_pending_requests\n raise RuntimeError" in str(log), log + + @deferred_f_from_coro_f + async def test_scheduler_enqueue_request_exception(self): + class TestScheduler(MemoryScheduler): + def enqueue_request(self, request): + raise RuntimeError + + class TestSpider(Spider): + name = "test" + start_urls = ["data:,"] + + def parse(self, response): + pass + + actual_dropped_urls = [] + + def track_dropped_url(request, spider): + actual_dropped_urls.append(request.url) + + settings = {"SCHEDULER": TestScheduler} + crawler = get_crawler(TestSpider, settings_dict=settings) + crawler.signals.connect(track_dropped_url, signals.request_dropped) + with LogCapture() as log: + await maybe_deferred_to_future(crawler.crawl()) + assert crawler.stats.get_value("finish_reason") == "finished" + expected_dropped_urls = ["data:,"] + assert actual_dropped_urls == expected_dropped_urls, ( + f"{actual_dropped_urls=} != {expected_dropped_urls=}" + ) + assert "in enqueue_request\n raise RuntimeError" in str(log), log + @deferred_f_from_coro_f async def test_scheduler_next_request_exception(self): class TestScheduler(MemoryScheduler):