Skip to main content

apify_client/clients/
request_queue.rs

1//! Client for a single request queue (`/v2/request-queues/{queueId}` and variants).
2
3use std::collections::HashSet;
4use std::time::Duration;
5
6use futures_util::stream::{self, StreamExt};
7use serde::Serialize;
8
9use crate::clients::base::{
10    delete_resource, delete_with_body, get_resource, get_resource_required, post_action,
11    post_with_body, update_resource, ResourceContext,
12};
13use crate::common::{encode_path_segment, QueryParams};
14use crate::error::{ApifyClientError, ApifyClientResult};
15use crate::http_client::{sleep_public, HttpClient, HttpMethod, HttpRequest};
16use crate::models::{
17    BatchRequestsOperationResult, LockedRequestQueueHead, RequestLockInfo, RequestQueue,
18    RequestQueueHead, RequestQueueOperationInfo, RequestQueueRequest, RequestQueueRequestsPage,
19    UnlockRequestsResult, UnprocessedRequest,
20};
21
22/// Maximum number of requests the API accepts in a single `requests/batch` call. Larger
23/// `batch_add_requests` inputs are split into chunks of at most this size (matching the
24/// reference client's `REQUEST_QUEUE_MAX_REQUESTS_PER_BATCH_OPERATION`); `batch_delete_requests`
25/// does not auto-chunk (matching the reference client) and instead rejects larger inputs.
26const MAX_REQUESTS_PER_BATCH_OPERATION: usize = 25;
27/// Maximum accepted size (bytes) of a request body, mirroring the platform-wide
28/// `MAX_PAYLOAD_SIZE_BYTES` (9 MiB) that the reference client chunks `batch_add_requests` calls
29/// against, on top of the per-call request-count limit.
30const MAX_PAYLOAD_SIZE_BYTES: usize = 9_437_184;
31/// Fraction of [`MAX_PAYLOAD_SIZE_BYTES`] held back as a safety margin (0.01%), matching the
32/// reference client's `SAFETY_BUFFER_PERCENT`.
33const SAFETY_BUFFER_PERCENT: f64 = 0.0001;
34/// Default number of batch-add API calls [`RequestQueueClient::batch_add_requests`] keeps in
35/// flight at once, matching the reference client's `DEFAULT_PARALLEL_BATCH_ADD_REQUESTS`.
36const DEFAULT_MAX_PARALLEL_BATCH_ADD_REQUESTS: usize = 5;
37/// Default number of retry attempts for requests a batch-add call reports as `unprocessed`
38/// (typically rate-limited), matching the reference client's
39/// `DEFAULT_UNPROCESSED_RETRIES_BATCH_ADD_REQUESTS`.
40const DEFAULT_MAX_UNPROCESSED_REQUESTS_RETRIES: u32 = 3;
41/// Default minimum delay before the first unprocessed-request retry; doubles (with jitter) on
42/// each subsequent retry, matching the reference client's
43/// `DEFAULT_MIN_DELAY_BETWEEN_UNPROCESSED_REQUESTS_RETRIES_MILLIS`.
44const DEFAULT_MIN_DELAY_BETWEEN_UNPROCESSED_REQUESTS_RETRIES: Duration = Duration::from_millis(500);
45
46/// Options for [`RequestQueueClient::batch_add_requests`].
47///
48/// Mirrors the reference client's retrying, chunked, parallel `batchAddRequests`: large inputs
49/// are split by count (max [`MAX_REQUESTS_PER_BATCH_OPERATION`]) and by JSON byte size (max
50/// [`MAX_PAYLOAD_SIZE_BYTES`], minus a safety margin), chunks are sent with up to `max_parallel`
51/// requests in flight at once, and any request an API call reports as `unprocessed` (typically
52/// due to rate limiting) is retried with exponential backoff.
53#[derive(Debug, Default, Clone)]
54pub struct BatchAddRequestsOptions {
55    /// If `true`, adds all requests to the front of the queue.
56    pub forefront: bool,
57    /// Maximum retries for requests reported as `unprocessed`. Defaults to
58    /// [`DEFAULT_MAX_UNPROCESSED_REQUESTS_RETRIES`] (3) when `None`.
59    pub max_unprocessed_requests_retries: Option<u32>,
60    /// Maximum number of chunk-add API calls in flight at once. Defaults to
61    /// [`DEFAULT_MAX_PARALLEL_BATCH_ADD_REQUESTS`] (5) when `None`.
62    pub max_parallel: Option<usize>,
63    /// Minimum delay before the first unprocessed-request retry (doubles, with jitter, on each
64    /// subsequent retry). Defaults to
65    /// [`DEFAULT_MIN_DELAY_BETWEEN_UNPROCESSED_REQUESTS_RETRIES`] (500ms) when `None`.
66    pub min_delay_between_unprocessed_requests_retries: Option<Duration>,
67}
68
69/// Returns the key used to correlate a request across the batch-add retry loop: its explicit
70/// `unique_key` if set, otherwise its `url` — matching the API's own fallback (a request added
71/// without a `unique_key` is deduplicated by its raw `url`).
72fn dedup_key(request: &RequestQueueRequest) -> &str {
73    request.unique_key.as_deref().unwrap_or(&request.url)
74}
75
76/// Returns the JSON-serialized byte length of `value`.
77fn json_byte_len<T: Serialize + ?Sized>(value: &T) -> ApifyClientResult<usize> {
78    Ok(serde_json::to_vec(value)?.len())
79}
80
81/// Slices `requests` down to a byte-limited prefix, mirroring the reference client's
82/// `sliceArrayByByteLength`: if the whole slice already fits under `max_bytes` it is returned
83/// unchanged; otherwise items are accumulated one at a time until the next one would exceed the
84/// budget. `start_index` is only used to name the offending item in the error message, so it
85/// should be the slice's absolute position within the caller's full input.
86///
87/// The first item is always included regardless of size (once its own size has been checked
88/// against `max_bytes`), guaranteeing a non-empty result for a non-empty input — unlike the
89/// reference implementation, which can return an empty slice (and loop forever) when a single
90/// item's size leaves no room under `max_bytes` for even itself plus the array wrapper.
91///
92/// Returns [`ApifyClientError::InvalidArgument`] if a single request's JSON exceeds `max_bytes`
93/// on its own (mirroring the reference client's thrown error).
94fn slice_requests_by_byte_length(
95    requests: &[RequestQueueRequest],
96    max_bytes: usize,
97    start_index: usize,
98) -> ApifyClientResult<Vec<RequestQueueRequest>> {
99    if json_byte_len(requests)? < max_bytes {
100        return Ok(requests.to_vec());
101    }
102    let mut out = Vec::new();
103    let mut byte_length = 2usize; // 2 bytes for the empty array `[]`.
104    for (offset, request) in requests.iter().enumerate() {
105        let item_bytes = json_byte_len(request)?;
106        if item_bytes > max_bytes {
107            return Err(ApifyClientError::InvalidArgument(format!(
108                "RequestQueueClient::batch_add_requests: the request at index {} exceeds the \
109                 maximum allowed size ({max_bytes} bytes)",
110                start_index + offset
111            )));
112        }
113        if !out.is_empty() && byte_length + item_bytes >= max_bytes {
114            break;
115        }
116        byte_length += item_bytes;
117        out.push(request.clone());
118    }
119    Ok(out)
120}
121
122/// Splits `requests` into chunks that each satisfy both the per-call count limit
123/// ([`MAX_REQUESTS_PER_BATCH_OPERATION`]) and the payload byte-size limit
124/// (`max_bytes`), mirroring the reference client's chunking loop in `batchAddRequests`.
125fn chunk_requests_for_batch_add(
126    requests: &[RequestQueueRequest],
127    max_bytes: usize,
128) -> ApifyClientResult<Vec<Vec<RequestQueueRequest>>> {
129    let mut chunks = Vec::new();
130    let mut i = 0;
131    while i < requests.len() {
132        let group_end = (i + MAX_REQUESTS_PER_BATCH_OPERATION).min(requests.len());
133        let chunk = slice_requests_by_byte_length(&requests[i..group_end], max_bytes, i)?;
134        // `slice_requests_by_byte_length` always returns at least one item for a non-empty
135        // input, so this advances on every iteration.
136        i += chunk.len();
137        chunks.push(chunk);
138    }
139    Ok(chunks)
140}
141
142/// Options for [`RequestQueueClient::list_requests`].
143///
144/// Covers the spec query parameters of `GET /v2/request-queues/{queueId}/requests`.
145#[derive(Debug, Default, Clone)]
146pub struct ListRequestsOptions {
147    /// Maximum number of requests to return.
148    pub limit: Option<i64>,
149    /// Start listing after this request ID (exclusive).
150    pub exclusive_start_id: Option<String>,
151    /// Opaque pagination cursor returned by a previous call.
152    pub cursor: Option<String>,
153    /// Restrict the returned requests to the given states. The spec defines this as an array of
154    /// the enum values `"locked"` and `"pending"`; multiple values are sent comma-joined (matching
155    /// the JS reference, which serializes `filter: Array<'locked' | 'pending'>` via `join(',')`).
156    pub filter: Option<Vec<String>>,
157}
158
159/// Client for a specific request queue.
160#[derive(Debug, Clone)]
161pub struct RequestQueueClient {
162    ctx: ResourceContext,
163    client_key: Option<String>,
164}
165
166impl RequestQueueClient {
167    pub(crate) fn new(http: HttpClient, base_url: &str, resource_path: &str, id: &str) -> Self {
168        Self {
169            ctx: ResourceContext::single(http, base_url, resource_path, id),
170            client_key: None,
171        }
172    }
173
174    /// Creates an RQ client for a run's default queue (nested path, no ID).
175    pub(crate) fn nested(http: HttpClient, base_url: &str, sub_path: &str) -> Self {
176        Self {
177            ctx: ResourceContext::collection(http, base_url, sub_path),
178            client_key: None,
179        }
180    }
181
182    /// Sets the `clientKey` used to identify this client across requests (for locking).
183    pub fn with_client_key(mut self, client_key: impl Into<String>) -> Self {
184        self.client_key = Some(client_key.into());
185        self
186    }
187
188    fn base_params(&self) -> QueryParams {
189        let mut params = QueryParams::new();
190        params.add_str("clientKey", self.client_key.clone());
191        params
192    }
193
194    /// Fetches the queue metadata, or `None` if it does not exist.
195    pub async fn get(&self) -> ApifyClientResult<Option<RequestQueue>> {
196        get_resource(&self.ctx, None, &QueryParams::new()).await
197    }
198
199    /// Updates the queue metadata (e.g. `name`, `title`).
200    pub async fn update<T: Serialize>(&self, new_fields: &T) -> ApifyClientResult<RequestQueue> {
201        update_resource(&self.ctx, None, new_fields).await
202    }
203
204    /// Deletes the queue.
205    pub async fn delete(&self) -> ApifyClientResult<()> {
206        delete_resource(&self.ctx, None).await
207    }
208
209    /// Lists requests from the head of the queue (without locking them).
210    pub async fn list_head(&self, limit: Option<i64>) -> ApifyClientResult<RequestQueueHead> {
211        let mut params = self.base_params();
212        params.add_int("limit", limit);
213        get_resource_required(&self.ctx, Some("head"), &params).await
214    }
215
216    /// Adds a single request to the queue. If `forefront` is true, adds it to the front.
217    pub async fn add_request(
218        &self,
219        request: &RequestQueueRequest,
220        forefront: bool,
221    ) -> ApifyClientResult<RequestQueueOperationInfo> {
222        let mut params = self.base_params();
223        params.add_bool("forefront", Some(forefront));
224        let body = serde_json::to_vec(request)?;
225        post_with_body(
226            &self.ctx,
227            Some("requests"),
228            &params,
229            Some(body),
230            "application/json",
231        )
232        .await
233    }
234
235    /// Gets a request by ID, or `None` if it does not exist.
236    pub async fn get_request(&self, id: &str) -> ApifyClientResult<Option<RequestQueueRequest>> {
237        get_resource(
238            &self.ctx,
239            Some(&format!("requests/{}", encode_path_segment(id))),
240            &self.base_params(),
241        )
242        .await
243    }
244
245    /// Updates a request (which must include its `id`).
246    pub async fn update_request(
247        &self,
248        request: &RequestQueueRequest,
249        forefront: bool,
250    ) -> ApifyClientResult<RequestQueueOperationInfo> {
251        let id = request.id.clone().ok_or_else(|| {
252            crate::error::ApifyClientError::InvalidArgument(
253                "request.id is required to update a request".to_string(),
254            )
255        })?;
256        let mut params = self.base_params();
257        params.add_bool("forefront", Some(forefront));
258        let url = params.apply_to_url(
259            &self
260                .ctx
261                .url(Some(&format!("requests/{}", encode_path_segment(&id)))),
262        );
263        let body = serde_json::to_vec(request)?;
264        let mut headers = std::collections::HashMap::new();
265        headers.insert("Content-Type".to_string(), "application/json".to_string());
266        let response = self
267            .ctx
268            .http
269            .call(HttpRequest {
270                method: HttpMethod::Put,
271                url,
272                headers,
273                body: Some(body),
274                timeout: crate::clients::base::DEFAULT_REQUEST_TIMEOUT,
275            })
276            .await?;
277        crate::common::parse_data_envelope(&response.body)
278    }
279
280    /// Deletes a request by ID.
281    pub async fn delete_request(&self, id: &str) -> ApifyClientResult<()> {
282        let params = self.base_params();
283        let url = params.apply_to_url(
284            &self
285                .ctx
286                .url(Some(&format!("requests/{}", encode_path_segment(id)))),
287        );
288        self.ctx
289            .http
290            .call(HttpRequest {
291                method: HttpMethod::Delete,
292                url,
293                headers: Default::default(),
294                body: None,
295                timeout: crate::clients::base::DEFAULT_REQUEST_TIMEOUT,
296            })
297            .await?;
298        Ok(())
299    }
300
301    /// Lists and locks requests from the head of the queue for `lock_secs` seconds.
302    pub async fn list_and_lock_head(
303        &self,
304        lock_secs: i64,
305        limit: Option<i64>,
306    ) -> ApifyClientResult<LockedRequestQueueHead> {
307        let mut params = self.base_params();
308        params
309            .add_int("lockSecs", Some(lock_secs))
310            .add_int("limit", limit);
311        post_action(&self.ctx, Some("head/lock"), &params, None, None).await
312    }
313
314    /// Adds multiple requests to the queue in a single logical operation.
315    ///
316    /// This is significantly more efficient than calling [`add_request`](Self::add_request)
317    /// once per request, especially for large batches: the input is automatically split into
318    /// chunks that respect both the API's per-call request-count limit
319    /// ([`MAX_REQUESTS_PER_BATCH_OPERATION`]) and its request-body byte-size limit
320    /// ([`MAX_PAYLOAD_SIZE_BYTES`]), chunks are sent with up to `options.max_parallel` API calls
321    /// in flight at once, and any request an API call reports as `unprocessed` (typically due to
322    /// rate limiting) is retried with exponential backoff — matching the reference client's
323    /// `batchAddRequests`. Every request must be identifiable by [`RequestQueueRequest::unique_key`]
324    /// (or, if left unset, by `url`, the API's own fallback) so a retried request can be matched
325    /// back to the original input.
326    ///
327    /// Unlike most methods here, this does not propagate per-chunk API errors: a chunk that fails
328    /// even after retries has its requests reported in the result's `unprocessed_requests`
329    /// instead, so a batch add of many requests never fails outright over one bad chunk (matching
330    /// the reference client). It still returns [`ApifyClientError::InvalidArgument`] up front,
331    /// before any request is sent, for an empty `requests` (matching
332    /// [`batch_delete_requests`](Self::batch_delete_requests) and the reference client, which
333    /// validates both as non-empty) or if a single request's JSON is too large to ever fit in a
334    /// chunk.
335    pub async fn batch_add_requests(
336        &self,
337        requests: &[RequestQueueRequest],
338        options: BatchAddRequestsOptions,
339    ) -> ApifyClientResult<BatchRequestsOperationResult> {
340        if requests.is_empty() {
341            return Err(ApifyClientError::InvalidArgument(
342                "RequestQueueClient::batch_add_requests requires at least 1 request".to_string(),
343            ));
344        }
345        let max_parallel = options
346            .max_parallel
347            .unwrap_or(DEFAULT_MAX_PARALLEL_BATCH_ADD_REQUESTS)
348            .max(1);
349        let payload_limit_bytes = MAX_PAYLOAD_SIZE_BYTES
350            - (MAX_PAYLOAD_SIZE_BYTES as f64 * SAFETY_BUFFER_PERCENT).ceil() as usize;
351        let chunks = chunk_requests_for_batch_add(requests, payload_limit_bytes)?;
352
353        let merged = stream::iter(chunks)
354            .map(|chunk| {
355                let client = self.clone();
356                let options = options.clone();
357                async move {
358                    client
359                        .batch_add_requests_chunk_with_retries(chunk, options)
360                        .await
361                }
362            })
363            .buffer_unordered(max_parallel)
364            .fold(
365                BatchRequestsOperationResult::default(),
366                |mut acc, chunk_result| async move {
367                    acc.processed_requests
368                        .extend(chunk_result.processed_requests);
369                    acc.unprocessed_requests
370                        .extend(chunk_result.unprocessed_requests);
371                    acc
372                },
373            )
374            .await;
375        Ok(merged)
376    }
377
378    /// Sends one already-byte/count-limited chunk, retrying requests reported as `unprocessed`
379    /// with exponential backoff (matching the reference client's `_batchAddRequestsWithRetries`).
380    ///
381    /// Never returns an `Err`: a transport/API failure that survives `HttpClient`'s own retries
382    /// marks every request still outstanding in this chunk as unprocessed instead of propagating,
383    /// so a single bad chunk cannot fail the whole (possibly-parallel) `batch_add_requests` call.
384    async fn batch_add_requests_chunk_with_retries(
385        &self,
386        chunk: Vec<RequestQueueRequest>,
387        options: BatchAddRequestsOptions,
388    ) -> BatchRequestsOperationResult {
389        let max_retries = options
390            .max_unprocessed_requests_retries
391            .unwrap_or(DEFAULT_MAX_UNPROCESSED_REQUESTS_RETRIES);
392        let min_delay = options
393            .min_delay_between_unprocessed_requests_retries
394            .unwrap_or(DEFAULT_MIN_DELAY_BETWEEN_UNPROCESSED_REQUESTS_RETRIES);
395
396        let mut remaining = chunk;
397        let mut processed = Vec::new();
398
399        for attempt in 0..=max_retries {
400            match self
401                .batch_add_requests_raw(&remaining, options.forefront)
402                .await
403            {
404                Ok(result) => {
405                    let processed_keys: HashSet<&str> = result
406                        .processed_requests
407                        .iter()
408                        .filter_map(|p| p.unique_key.as_deref())
409                        .collect();
410                    remaining.retain(|r| !processed_keys.contains(dedup_key(r)));
411                    processed.extend(result.processed_requests);
412                    if remaining.is_empty() {
413                        return BatchRequestsOperationResult {
414                            processed_requests: processed,
415                            unprocessed_requests: Vec::new(),
416                        };
417                    }
418                    if attempt == max_retries {
419                        return BatchRequestsOperationResult {
420                            processed_requests: processed,
421                            unprocessed_requests: result.unprocessed_requests,
422                        };
423                    }
424                }
425                Err(_) => {
426                    // A hard failure (already retried by `HttpClient` for transient errors):
427                    // treat every request still outstanding in this chunk as unprocessed rather
428                    // than propagating, matching the reference client's "never throws" contract.
429                    return BatchRequestsOperationResult {
430                        processed_requests: processed,
431                        unprocessed_requests: remaining
432                            .iter()
433                            .map(|r| UnprocessedRequest {
434                                unique_key: dedup_key(r).to_string(),
435                                url: r.url.clone(),
436                                method: r.method.clone(),
437                            })
438                            .collect(),
439                    };
440                }
441            }
442            // Exponential backoff with jitter before the next retry, matching the reference
443            // client's `(1 + random()) * 2^attempt * minDelay` (see `randomized_delay`, which
444            // returns a value in `[base, 2*base)`, i.e. `(1 + random()) * base`).
445            let backoff = min_delay.saturating_mul(2u32.saturating_pow(attempt));
446            sleep_public(crate::http_client::randomized_delay(backoff)).await;
447        }
448        // Unreachable: the loop above always returns on its last iteration (`attempt ==
449        // max_retries` is handled inside the `Ok` arm, and `Err` returns unconditionally). Kept
450        // as a safe fallback rather than `unreachable!()` so a future refactor of the loop bounds
451        // fails safe (reporting the batch unprocessed) instead of panicking.
452        BatchRequestsOperationResult {
453            processed_requests: processed,
454            unprocessed_requests: remaining
455                .iter()
456                .map(|r| UnprocessedRequest {
457                    unique_key: dedup_key(r).to_string(),
458                    url: r.url.clone(),
459                    method: r.method.clone(),
460                })
461                .collect(),
462        }
463    }
464
465    /// Sends a single `POST requests/batch` call (at most [`MAX_REQUESTS_PER_BATCH_OPERATION`]
466    /// requests, and within the byte-size budget already enforced by the caller).
467    async fn batch_add_requests_raw(
468        &self,
469        requests: &[RequestQueueRequest],
470        forefront: bool,
471    ) -> ApifyClientResult<BatchRequestsOperationResult> {
472        let mut params = self.base_params();
473        params.add_bool("forefront", Some(forefront));
474        let body = serde_json::to_vec(requests)?;
475        post_with_body(
476            &self.ctx,
477            Some("requests/batch"),
478            &params,
479            Some(body),
480            "application/json",
481        )
482        .await
483    }
484
485    /// Deletes multiple requests in a single batch operation.
486    ///
487    /// Unlike [`batch_add_requests`](Self::batch_add_requests), this does not auto-chunk: the API
488    /// accepts at most [`MAX_REQUESTS_PER_BATCH_OPERATION`] requests per call (matching the
489    /// reference client, which validates rather than chunks), so a larger `requests` returns
490    /// [`ApifyClientError::InvalidArgument`] before any request is sent.
491    pub async fn batch_delete_requests<T: Serialize>(
492        &self,
493        requests: &[T],
494    ) -> ApifyClientResult<BatchRequestsOperationResult> {
495        if requests.is_empty() || requests.len() > MAX_REQUESTS_PER_BATCH_OPERATION {
496            return Err(ApifyClientError::InvalidArgument(format!(
497                "RequestQueueClient::batch_delete_requests accepts between 1 and {MAX_REQUESTS_PER_BATCH_OPERATION} requests per call, got {}",
498                requests.len()
499            )));
500        }
501        delete_with_body(
502            &self.ctx,
503            Some("requests/batch"),
504            &self.base_params(),
505            &requests,
506        )
507        .await
508    }
509
510    /// Lists requests in the queue.
511    ///
512    /// Supports pagination via `limit`/`exclusive_start_id` and the spec's `cursor`/`filter`
513    /// parameters (see [`ListRequestsOptions`]).
514    pub async fn list_requests(
515        &self,
516        options: ListRequestsOptions,
517    ) -> ApifyClientResult<RequestQueueRequestsPage> {
518        let mut params = self.base_params();
519        params
520            .add_int("limit", options.limit)
521            .add_str("exclusiveStartId", options.exclusive_start_id)
522            .add_str("cursor", options.cursor)
523            .add_csv("filter", options.filter.as_deref());
524        get_resource_required(&self.ctx, Some("requests"), &params).await
525    }
526
527    /// Prolongs the lock on a request for another `lock_secs` seconds.
528    ///
529    /// If `forefront` is `true`, the request moves to the front of the queue when its lock
530    /// later expires.
531    pub async fn prolong_request_lock(
532        &self,
533        id: &str,
534        lock_secs: i64,
535        forefront: bool,
536    ) -> ApifyClientResult<RequestLockInfo> {
537        let mut params = self.base_params();
538        params
539            .add_int("lockSecs", Some(lock_secs))
540            .add_bool("forefront", Some(forefront));
541        let url = params.apply_to_url(
542            &self
543                .ctx
544                .url(Some(&format!("requests/{}/lock", encode_path_segment(id)))),
545        );
546        let response = self
547            .ctx
548            .http
549            .call(HttpRequest {
550                method: HttpMethod::Put,
551                url,
552                headers: Default::default(),
553                body: None,
554                timeout: crate::clients::base::MEDIUM_REQUEST_TIMEOUT,
555            })
556            .await?;
557        crate::common::parse_data_envelope(&response.body)
558    }
559
560    /// Releases the lock on a request so other clients can process it.
561    ///
562    /// If `forefront` is `true`, the request moves to the front of the queue.
563    pub async fn delete_request_lock(&self, id: &str, forefront: bool) -> ApifyClientResult<()> {
564        let mut params = self.base_params();
565        params.add_bool("forefront", Some(forefront));
566        let url = params.apply_to_url(
567            &self
568                .ctx
569                .url(Some(&format!("requests/{}/lock", encode_path_segment(id)))),
570        );
571        self.ctx
572            .http
573            .call(HttpRequest {
574                method: HttpMethod::Delete,
575                url,
576                headers: Default::default(),
577                body: None,
578                timeout: crate::clients::base::SMALL_REQUEST_TIMEOUT,
579            })
580            .await?;
581        Ok(())
582    }
583
584    /// Lazily paginates over all requests in the queue, fetching pages on demand.
585    ///
586    /// Returns a [`RequestQueueRequestsIterator`]; call its `next()` to get one request at a
587    /// time. Pagination uses the API's opaque `nextCursor` token: the first page may be
588    /// anchored with `exclusiveStartId`, but every subsequent page is fetched with `cursor`
589    /// (matching the JS reference). `cursor` and `exclusiveStartId` are mutually exclusive.
590    pub fn paginate_requests(&self, page_limit: Option<i64>) -> RequestQueueRequestsIterator {
591        RequestQueueRequestsIterator {
592            client: self.clone(),
593            page_limit,
594            buffer: std::collections::VecDeque::new(),
595            next_cursor: None,
596            exhausted: false,
597        }
598    }
599
600    /// Unlocks all requests currently locked by this client (identified by `client_key`).
601    pub async fn unlock_requests(&self) -> ApifyClientResult<UnlockRequestsResult> {
602        post_action(
603            &self.ctx,
604            Some("requests/unlock"),
605            &self.base_params(),
606            None,
607            None,
608        )
609        .await
610    }
611}
612
613/// A lazy, page-fetching iterator over the requests in a queue.
614///
615/// Created by [`RequestQueueClient::paginate_requests`]. Each call to [`next`](Self::next)
616/// returns the next request, fetching another page from the API when the local buffer is
617/// exhausted, until all requests have been yielded.
618pub struct RequestQueueRequestsIterator {
619    client: RequestQueueClient,
620    page_limit: Option<i64>,
621    buffer: std::collections::VecDeque<RequestQueueRequest>,
622    /// Opaque pagination token returned by the previous page, fed back as `cursor`.
623    next_cursor: Option<String>,
624    exhausted: bool,
625}
626
627impl RequestQueueRequestsIterator {
628    /// Returns the next request, or `None` when all requests have been yielded.
629    pub async fn next(&mut self) -> ApifyClientResult<Option<RequestQueueRequest>> {
630        if let Some(item) = self.buffer.pop_front() {
631            return Ok(Some(item));
632        }
633        if self.exhausted {
634            return Ok(None);
635        }
636
637        // The first page may be anchored by exclusiveStartId; every later page is fetched
638        // with the opaque `cursor` token (mutually exclusive with exclusiveStartId), matching
639        // the JS reference. Here we only ever paginate from the queue head, so the first page
640        // uses neither and subsequent pages use `cursor`.
641        let page = self
642            .client
643            .list_requests(ListRequestsOptions {
644                limit: self.page_limit,
645                cursor: self.next_cursor.clone(),
646                ..Default::default()
647            })
648            .await?;
649
650        if page.items.is_empty() {
651            self.exhausted = true;
652            return Ok(None);
653        }
654
655        // Advance the cursor; stop when the API stops returning one.
656        match page.next_cursor {
657            Some(cursor) if !cursor.is_empty() => self.next_cursor = Some(cursor),
658            _ => self.exhausted = true,
659        }
660
661        self.buffer.extend(page.items);
662        Ok(self.buffer.pop_front())
663    }
664}
665
666#[cfg(test)]
667mod batch_add_tests {
668    use super::{
669        chunk_requests_for_batch_add, dedup_key, slice_requests_by_byte_length,
670        MAX_REQUESTS_PER_BATCH_OPERATION,
671    };
672    use crate::models::RequestQueueRequest;
673
674    fn request(url: &str, unique_key: Option<&str>) -> RequestQueueRequest {
675        RequestQueueRequest {
676            id: None,
677            url: url.to_string(),
678            unique_key: unique_key.map(str::to_string),
679            method: None,
680            user_data: None,
681            extra: Default::default(),
682        }
683    }
684
685    /// A request without an explicit `unique_key` is correlated by `url`, matching the API's own
686    /// deduplication fallback.
687    #[test]
688    fn dedup_key_falls_back_to_url() {
689        let with_key = request("https://example.com", Some("k1"));
690        assert_eq!(dedup_key(&with_key), "k1");
691
692        let without_key = request("https://example.com/no-key", None);
693        assert_eq!(dedup_key(&without_key), "https://example.com/no-key");
694    }
695
696    /// A slice that already fits under the byte budget is returned unchanged.
697    #[test]
698    fn byte_slice_returns_everything_when_under_budget() {
699        let requests: Vec<_> = (0..5)
700            .map(|i| request(&format!("https://example.com/{i}"), None))
701            .collect();
702        let sliced = slice_requests_by_byte_length(&requests, 1_000_000, 0).unwrap();
703        assert_eq!(sliced.len(), 5);
704    }
705
706    /// When the whole slice exceeds the byte budget, only a byte-limited prefix is taken — but
707    /// never an empty one, even if the very first item alone leaves no room for a second.
708    #[test]
709    fn byte_slice_takes_a_limited_prefix() {
710        let requests: Vec<_> = (0..10)
711            .map(|i| request(&format!("https://example.com/{i}"), None))
712            .collect();
713        // Each request serializes to roughly 30 bytes; a budget of 50 fits one comfortably but
714        // never two.
715        let sliced = slice_requests_by_byte_length(&requests, 50, 0).unwrap();
716        assert_eq!(
717            sliced.len(),
718            1,
719            "budget of 50 bytes should admit exactly one ~30-byte request"
720        );
721    }
722
723    /// A single request whose own JSON exceeds the byte budget is a hard error, not a silently
724    /// dropped item — the caller could never send it in any chunk.
725    #[test]
726    fn byte_slice_errors_on_oversized_single_request() {
727        let huge_url = format!("https://example.com/{}", "x".repeat(1000));
728        let requests = vec![request(&huge_url, None)];
729        let err = slice_requests_by_byte_length(&requests, 100, 3).unwrap_err();
730        let message = err.to_string();
731        assert!(
732            message.contains("index 3"),
733            "error should name the absolute index of the oversized request: {message}"
734        );
735    }
736
737    /// Chunking respects the per-call count cap even when every request is tiny (byte budget is
738    /// never the limiting factor).
739    #[test]
740    fn chunking_splits_by_count_when_bytes_are_plentiful() {
741        let requests: Vec<_> = (0..(MAX_REQUESTS_PER_BATCH_OPERATION * 2 + 3))
742            .map(|i| request(&format!("https://example.com/{i}"), None))
743            .collect();
744        let chunks = chunk_requests_for_batch_add(&requests, 10_000_000).unwrap();
745        let sizes: Vec<usize> = chunks.iter().map(Vec::len).collect();
746        assert_eq!(
747            sizes,
748            vec![
749                MAX_REQUESTS_PER_BATCH_OPERATION,
750                MAX_REQUESTS_PER_BATCH_OPERATION,
751                3
752            ]
753        );
754        let total: usize = sizes.iter().sum();
755        assert_eq!(total, requests.len());
756    }
757
758    /// Chunking also respects the byte budget, producing more (smaller) chunks than the count
759    /// cap alone would when requests are large.
760    #[test]
761    fn chunking_splits_by_byte_budget_when_tighter_than_count_cap() {
762        let requests: Vec<_> = (0..6)
763            .map(|i| request(&format!("https://example.com/{i}"), None))
764            .collect();
765        // ~30 bytes/request; a 100-byte budget forces multiple chunks well under the 25-item cap.
766        let chunks = chunk_requests_for_batch_add(&requests, 100).unwrap();
767        assert!(
768            chunks.len() > 1,
769            "a tight byte budget must force more than one chunk, got {}",
770            chunks.len()
771        );
772        let total: usize = chunks.iter().map(Vec::len).sum();
773        assert_eq!(
774            total,
775            requests.len(),
776            "every request must end up in exactly one chunk"
777        );
778    }
779}