Improve exception handling for Scheduler’s has_pending_requests() and enqueue_request()

This commit is contained in:
Adrián Chaves 2025-03-15 20:50:21 +01:00
parent c8bbe4d968
commit a5dc229508
2 changed files with 88 additions and 3 deletions

View File

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

View File

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