Repository navigation
Expand file tree
/
Copy path_request_queue_client.py
More file actions
144 lines (108 loc) · 5.21 KB
/
Copy path_request_queue_client.py
File metadata and controls
144 lines (108 loc) · 5.21 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
from __future__ import annotations
from abc import ABC, abstractmethod
from typing import TYPE_CHECKING
if TYPE_CHECKING:
from collections.abc import Sequence
from crawlee import Request
from crawlee.storage_clients.models import AddRequestsResponse, ProcessedRequest, RequestQueueMetadata
class RequestQueueClient(ABC):
"""An abstract class for request queue resource clients.
These clients are specific to the type of resource they manage and operate under a designated storage
client, like a memory storage client.
"""
@abstractmethod
async def get_metadata(self) -> RequestQueueMetadata:
"""Get the metadata of the request queue."""
@abstractmethod
async def drop(self) -> None:
"""Drop the whole request queue and remove all its values.
The backend method for the `RequestQueue.drop` call.
"""
@abstractmethod
async def purge(self) -> None:
"""Purge all items from the request queue.
The backend method for the `RequestQueue.purge` call.
"""
@abstractmethod
async def add_batch_of_requests(
self,
requests: Sequence[Request],
*,
forefront: bool = False,
) -> AddRequestsResponse:
"""Add batch of requests to the queue.
This method adds a batch of requests to the queue. Each request is processed based on its uniqueness
(determined by `unique_key`). Duplicates will be identified but not re-added to the queue.
Args:
requests: The collection of requests to add to the queue.
forefront: Whether to put the added requests at the beginning (True) or the end (False) of the queue.
When True, the requests will be processed sooner than previously added requests.
Returns:
A response object containing information about which requests were successfully
processed and which failed (if any).
"""
@abstractmethod
async def get_request(self, unique_key: str) -> Request | None:
"""Retrieve a request from the queue.
Args:
unique_key: Unique key of the request to retrieve.
Returns:
The retrieved request, or None, if it did not exist.
"""
@abstractmethod
async def fetch_next_request(self) -> Request | None:
"""Return the next request in the queue to be processed.
Once you successfully finish processing of the request, you need to call `RequestQueue.mark_request_as_handled`
to mark the request as handled in the queue. If there was some error in processing the request, call
`RequestQueue.reclaim_request` instead, so that the queue will give the request to some other consumer
in another call to the `fetch_next_request` method.
Note that the `None` return value does not mean the queue processing finished, it means there are currently
no pending requests. To check whether all requests in queue were finished, use `RequestQueue.is_finished`
instead.
Returns:
The request or `None` if there are no more pending requests.
"""
@abstractmethod
async def mark_request_as_handled(self, request: Request) -> ProcessedRequest | None:
"""Mark a request as handled after successful processing.
Handled requests will never again be returned by the `RequestQueue.fetch_next_request` method.
Args:
request: The request to mark as handled.
Returns:
Information about the queue operation. `None` if the given request was not in progress.
"""
@abstractmethod
async def reclaim_request(
self,
request: Request,
*,
forefront: bool = False,
) -> ProcessedRequest | None:
"""Reclaim a failed request back to the queue.
The request will be returned for processing later again by another call to `RequestQueue.fetch_next_request`.
Args:
request: The request to return to the queue.
forefront: Whether to add the request to the head or the end of the queue.
Returns:
Information about the queue operation. `None` if the given request was not in progress.
"""
@abstractmethod
async def is_empty(self) -> bool:
"""Check if the request queue is empty. That means there are no requests available to fetch.
Returns:
True if the request queue is empty, False otherwise.
"""
async def is_finished(self) -> bool:
"""Check if the request queue is finished.
A finished queue is empty and has no requests currently being processed.
Warning:
This default only checks `is_empty`, which reports whether requests are available to fetch, not
whether requests are still being processed. It can therefore return `True` while requests are in
progress. Subclasses that track in-progress requests should override this method; all built-in
clients already do.
Returns:
True if the request queue is finished, False otherwise.
"""
# TODO: Make this method abstract.
# https://github.com/apify/crawlee-python/issues/1985
return await self.is_empty()