From f5d41283c304e7c1531cbc6bac8354f76c92b2b8 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Fri, 31 Jul 2026 04:57:40 +0200 Subject: [PATCH 1/3] Align peek and pop --- scrapy/pqueues.py | 76 +++++++++++++------------ tests/test_pqueues.py | 125 ++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 165 insertions(+), 36 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 41411ceaf..abbcbb6da 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -151,7 +151,10 @@ class ScrapyPriorityQueue: else: q.close() - self.curprio = min(startprios) + # A recorded priority may have no queue to restore, e.g. if it only + # ever held a request that failed to serialize, so curprio comes from + # the queues that do exist and not from startprios. + self._update_curprio() def qfactory(self, key: int) -> QueueProtocol: return build_from_crawler( @@ -188,41 +191,37 @@ class ScrapyPriorityQueue: def pop(self) -> Request | None: while self.curprio is not None: - try: - q = self.queues[self.curprio] - except KeyError: - pass - else: + for queues in (self.queues, self._start_queues): + q = queues.get(self.curprio) + # An empty queue can linger at a priority when a push failed + # after creating it, e.g. on a serialization error. Popping + # from it would return None and hide the request that the other + # dict may hold at the same priority. + if not q: + continue m = q.pop() if not q: - del self.queues[self.curprio] - q.close() - if not self._start_queues: - self._update_curprio() - return m - if self._start_queues: - try: - q = self._start_queues[self.curprio] - except KeyError: + # The other dict may have no queue at this priority either, + # and a curprio that neither dict has would make peek() come + # up empty. self._update_curprio() - else: - m = q.pop() - if not q: - del self._start_queues[self.curprio] - q.close() - self._update_curprio() - return m - else: - self._update_curprio() + return m + # Nothing to pop at this priority: refreshing drops the empty + # leftovers and moves on to the next priority. + self._update_curprio() return None def _update_curprio(self) -> None: - prios = { - p - for queues in (self.queues, self._start_queues) - for p, q in queues.items() - if q - } + # Keeping an empty queue would hold its storage open for nothing, and + # make close() record a priority with nothing to restore from it. + prios: set[int] = set() + for queues in (self.queues, self._start_queues): + for p, q in list(queues.items()): + if q: + prios.add(p) + else: + del queues[p] + q.close() self.curprio = min(prios) if prios else None def peek(self) -> Request | None: @@ -234,12 +233,17 @@ class ScrapyPriorityQueue: """ if self.curprio is None: return None - try: - queue = self._start_queues[self.curprio] - except KeyError: - queue = self.queues[self.curprio] - # Protocols can't declare optional members - return cast("Request", queue.peek()) # type: ignore[attr-defined] + # The dicts are walked in the same order as in pop(), which is what + # makes the returned request the one that pop() then returns. + for queues in (self.queues, self._start_queues): + queue = queues.get(self.curprio) + # Empty queues can linger at a priority (see pop()), where they + # would hide the request that the other dict may hold at the same + # priority. + if queue: + # Protocols can't declare optional members + return cast("Request", queue.peek()) # type: ignore[attr-defined] + return None def close(self) -> list[int]: active: set[int] = set() diff --git a/tests/test_pqueues.py b/tests/test_pqueues.py index 85fefd172..b0a804a8e 100644 --- a/tests/test_pqueues.py +++ b/tests/test_pqueues.py @@ -76,6 +76,131 @@ class TestPriorityQueue: assert queue.pop().url == req3.url assert not queue.close() + def test_peek_after_draining_a_higher_priority_queue(self): + """Draining the queue of the current priority while start requests + remain at a different one must not leave ``curprio`` pointing at a + priority that no queue has, which would make ``peek()`` come up empty + with requests still queued.""" + temp_dir = tempfile.mkdtemp() + queue = ScrapyPriorityQueue.from_crawler( + self.crawler, + FifoMemoryQueue, + temp_dir, + start_queue_cls=FifoMemoryQueue, + ) + start_request = Request( + "https://example.org/start", meta={"is_start_request": True} + ) + queue.push(start_request) + # A redirect of a start request, which REDIRECT_PRIORITY_ADJUST puts at + # a higher priority, i.e. in a separate, non-start queue. + queue.push(Request("https://example.org/redirect", priority=2)) + + assert queue.peek().url == "https://example.org/redirect" + assert queue.pop().url == "https://example.org/redirect" + assert queue.peek().url == start_request.url + assert queue.pop().url == start_request.url + assert queue.peek() is None + queue.close() + + def test_peek_agrees_with_pop_on_start_requests(self): + """A start request and a non-start request at the same priority sit in + separate queues, and ``peek()`` must report the one that ``pop()`` + returns, since a caller may peek to decide whether to pop.""" + temp_dir = tempfile.mkdtemp() + queue = ScrapyPriorityQueue.from_crawler( + self.crawler, + FifoMemoryQueue, + temp_dir, + start_queue_cls=FifoMemoryQueue, + ) + queue.push( + Request("https://example.org/start", meta={"is_start_request": True}) + ) + queue.push(Request("https://example.org/other")) + + while len(queue): + peeked = queue.peek() + assert queue.pop().url == peeked.url + queue.close() + + def test_peek_and_pop_skip_an_empty_queue_left_by_a_failed_push(self): + """A queue is created before the request is pushed into it, so a push + that fails (e.g. a serialization error) leaves an empty queue behind at + that priority. Neither ``peek()`` nor ``pop()`` may let it hide the + request that the other dict holds at the same priority.""" + temp_dir = tempfile.mkdtemp() + queue = ScrapyPriorityQueue.from_crawler( + self.crawler, + PickleFifoDiskQueue, + temp_dir, + start_queue_cls=PickleFifoDiskQueue, + ) + with pytest.raises(ValueError, match="is not an instance method"): + queue.push( + Request("https://example.org/lambda", callback=lambda response: None) + ) + assert queue.queues[0] is not None # the empty leftover + assert len(queue) == 0 + + start_request = Request( + "https://example.org/start", meta={"is_start_request": True} + ) + queue.push(start_request) + assert len(queue) == 1 + assert queue.peek().url == start_request.url + assert queue.pop().url == start_request.url + assert len(queue) == 0 + assert queue.peek() is None + assert queue.pop() is None + queue.close() + + def test_empty_queues_are_dropped_on_refresh(self): + """The empty leftover of a failed push is forgotten (and closed) the + next time the current priority is refreshed, rather than kept around + holding its storage open.""" + temp_dir = tempfile.mkdtemp() + queue = ScrapyPriorityQueue.from_crawler( + self.crawler, PickleFifoDiskQueue, temp_dir + ) + with pytest.raises(ValueError, match="is not an instance method"): + queue.push( + Request("https://example.org/lambda", callback=lambda response: None) + ) + assert set(queue.queues) == {0} # the empty leftover + # A request at a different priority, whose queue emptying is what + # triggers the refresh. + queue.push(Request("https://example.org/1", priority=1)) + assert queue.pop().url == "https://example.org/1" + assert queue.queues == {} + assert queue.curprio is None + assert queue.pop() is None + assert not queue.close() + + def test_init_prios_without_a_restorable_queue(self): + """A priority recorded on close may have nothing to restore, e.g. if + its only request failed to serialize. ``curprio`` must not point at it, + or ``peek()`` raises :exc:`KeyError` on the resumed crawl.""" + temp_dir = tempfile.mkdtemp() + queue = ScrapyPriorityQueue.from_crawler( + self.crawler, PickleFifoDiskQueue, temp_dir + ) + with pytest.raises(ValueError, match="is not an instance method"): + queue.push( + Request("https://example.org/lambda", callback=lambda response: None) + ) + startprios = queue.close() + assert startprios == [0] + + queue2 = ScrapyPriorityQueue.from_crawler( + self.crawler, PickleFifoDiskQueue, temp_dir, startprios + ) + assert len(queue2) == 0 + assert queue2.curprio is None + assert queue2.peek() is None + assert queue2.pop() is None + queue2.close() + def test_init_prios_with_start_queue(self): temp_dir = tempfile.mkdtemp() queue = ScrapyPriorityQueue.from_crawler( From 2f5821f5f6cd736f78870b03e51f89aa1583e042 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Fri, 31 Jul 2026 05:16:26 +0200 Subject: [PATCH 2/3] Skip min peek() tests --- tests/test_pqueues.py | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/tests/test_pqueues.py b/tests/test_pqueues.py index b0a804a8e..8ad10e1c6 100644 --- a/tests/test_pqueues.py +++ b/tests/test_pqueues.py @@ -81,6 +81,8 @@ class TestPriorityQueue: remain at a different one must not leave ``curprio`` pointing at a priority that no queue has, which would make ``peek()`` come up empty with requests still queued.""" + if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + pytest.skip("queuelib.queue.FifoMemoryQueue.peek is undefined") temp_dir = tempfile.mkdtemp() queue = ScrapyPriorityQueue.from_crawler( self.crawler, @@ -107,6 +109,8 @@ class TestPriorityQueue: """A start request and a non-start request at the same priority sit in separate queues, and ``peek()`` must report the one that ``pop()`` returns, since a caller may peek to decide whether to pop.""" + if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + pytest.skip("queuelib.queue.FifoMemoryQueue.peek is undefined") temp_dir = tempfile.mkdtemp() queue = ScrapyPriorityQueue.from_crawler( self.crawler, @@ -129,6 +133,8 @@ class TestPriorityQueue: that fails (e.g. a serialization error) leaves an empty queue behind at that priority. Neither ``peek()`` nor ``pop()`` may let it hide the request that the other dict holds at the same priority.""" + if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + pytest.skip("queuelib.queue.FifoMemoryQueue.peek is undefined") temp_dir = tempfile.mkdtemp() queue = ScrapyPriorityQueue.from_crawler( self.crawler, From 2efa0b6b439cde778d8af4a611c1070860b84d29 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Fri, 31 Jul 2026 15:25:08 +0200 Subject: [PATCH 3/3] Improve test coverage --- scrapy/pqueues.py | 53 ++++++++++++++++++------------------------- tests/test_pqueues.py | 34 +++++++++++++++++++++++++++ 2 files changed, 56 insertions(+), 31 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index abbcbb6da..3bdc7857d 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -186,30 +186,27 @@ class ScrapyPriorityQueue: self.queues[priority] = self.qfactory(priority) q = self.queues[priority] q.push(request) # this may fail (eg. serialization error) - if self.curprio is None or priority < self.curprio: + # A queue class may drop a request instead of storing it, and a + # priority with an empty queue must not become the current one, since + # pop() and peek() expect to find a request there. + if q and (self.curprio is None or priority < self.curprio): self.curprio = priority + def _curqueue(self) -> QueueProtocol: + assert self.curprio is not None + # Whichever dict holds a non-empty queue at the current priority. The + # other one may hold an empty queue left behind by a failed push, e.g. + # on a serialization error, which would hide the next request. + return self.queues.get(self.curprio) or self._start_queues[self.curprio] + def pop(self) -> Request | None: - while self.curprio is not None: - for queues in (self.queues, self._start_queues): - q = queues.get(self.curprio) - # An empty queue can linger at a priority when a push failed - # after creating it, e.g. on a serialization error. Popping - # from it would return None and hide the request that the other - # dict may hold at the same priority. - if not q: - continue - m = q.pop() - if not q: - # The other dict may have no queue at this priority either, - # and a curprio that neither dict has would make peek() come - # up empty. - self._update_curprio() - return m - # Nothing to pop at this priority: refreshing drops the empty - # leftovers and moves on to the next priority. + if self.curprio is None: + return None + q = self._curqueue() + request = q.pop() + if not q: self._update_curprio() - return None + return request def _update_curprio(self) -> None: # Keeping an empty queue would hold its storage open for nothing, and @@ -233,17 +230,11 @@ class ScrapyPriorityQueue: """ if self.curprio is None: return None - # The dicts are walked in the same order as in pop(), which is what - # makes the returned request the one that pop() then returns. - for queues in (self.queues, self._start_queues): - queue = queues.get(self.curprio) - # Empty queues can linger at a priority (see pop()), where they - # would hide the request that the other dict may hold at the same - # priority. - if queue: - # Protocols can't declare optional members - return cast("Request", queue.peek()) # type: ignore[attr-defined] - return None + # Picking the queue the same way as pop() is what makes the returned + # request the one that pop() then returns. + queue = self._curqueue() + # Protocols can't declare optional members + return cast("Request", queue.peek()) # type: ignore[attr-defined] def close(self) -> list[int]: active: set[int] = set() diff --git a/tests/test_pqueues.py b/tests/test_pqueues.py index 8ad10e1c6..a30f8f3b0 100644 --- a/tests/test_pqueues.py +++ b/tests/test_pqueues.py @@ -14,6 +14,15 @@ from scrapy.utils.test import get_crawler from tests.utils.downloader import MockDownloader +class DroppingFifoMemoryQueue(FifoMemoryQueue): # type: ignore[valid-type,misc] + """Queue that drops requests marked with ``drop`` in their metadata, the + way a custom queue class could drop requests it does not want to store.""" + + def push(self, request): + if not request.meta.get("drop"): + super().push(request) + + class TestPriorityQueue: def setup_method(self): self.crawler = get_crawler(Spider) @@ -183,6 +192,31 @@ class TestPriorityQueue: assert queue.pop() is None assert not queue.close() + def test_push_into_a_queue_that_drops_the_request(self): + """A queue class may drop a request instead of storing it, and that + priority must not become the current one, or ``peek()`` and ``pop()`` + would come up empty with requests still queued.""" + if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + pytest.skip("queuelib.queue.FifoMemoryQueue.peek is undefined") + temp_dir = tempfile.mkdtemp() + queue = ScrapyPriorityQueue.from_crawler( + self.crawler, DroppingFifoMemoryQueue, temp_dir + ) + # A higher priority than the request below, so that it would become the + # current priority if it were stored. + queue.push( + Request("https://example.org/dropped", priority=1, meta={"drop": True}) + ) + kept = Request("https://example.org/kept") + queue.push(kept) + + assert len(queue) == 1 + assert queue.peek().url == kept.url + assert queue.pop().url == kept.url + assert queue.peek() is None + assert queue.pop() is None + assert not queue.close() + def test_init_prios_without_a_restorable_queue(self): """A priority recorded on close may have nothing to restore, e.g. if its only request failed to serialize. ``curprio`` must not point at it,