1use 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
22const MAX_REQUESTS_PER_BATCH_OPERATION: usize = 25;
27const MAX_PAYLOAD_SIZE_BYTES: usize = 9_437_184;
31const SAFETY_BUFFER_PERCENT: f64 = 0.0001;
34const DEFAULT_MAX_PARALLEL_BATCH_ADD_REQUESTS: usize = 5;
37const DEFAULT_MAX_UNPROCESSED_REQUESTS_RETRIES: u32 = 3;
41const DEFAULT_MIN_DELAY_BETWEEN_UNPROCESSED_REQUESTS_RETRIES: Duration = Duration::from_millis(500);
45
46#[derive(Debug, Default, Clone)]
54pub struct BatchAddRequestsOptions {
55 pub forefront: bool,
57 pub max_unprocessed_requests_retries: Option<u32>,
60 pub max_parallel: Option<usize>,
63 pub min_delay_between_unprocessed_requests_retries: Option<Duration>,
67}
68
69fn dedup_key(request: &RequestQueueRequest) -> &str {
73 request.unique_key.as_deref().unwrap_or(&request.url)
74}
75
76fn json_byte_len<T: Serialize + ?Sized>(value: &T) -> ApifyClientResult<usize> {
78 Ok(serde_json::to_vec(value)?.len())
79}
80
81fn 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; 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
122fn 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 i += chunk.len();
137 chunks.push(chunk);
138 }
139 Ok(chunks)
140}
141
142#[derive(Debug, Default, Clone)]
146pub struct ListRequestsOptions {
147 pub limit: Option<i64>,
149 pub exclusive_start_id: Option<String>,
151 pub cursor: Option<String>,
153 pub filter: Option<Vec<String>>,
157}
158
159#[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 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 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 pub async fn get(&self) -> ApifyClientResult<Option<RequestQueue>> {
196 get_resource(&self.ctx, None, &QueryParams::new()).await
197 }
198
199 pub async fn update<T: Serialize>(&self, new_fields: &T) -> ApifyClientResult<RequestQueue> {
201 update_resource(&self.ctx, None, new_fields).await
202 }
203
204 pub async fn delete(&self) -> ApifyClientResult<()> {
206 delete_resource(&self.ctx, None).await
207 }
208
209 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"), ¶ms).await
214 }
215
216 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 ¶ms,
229 Some(body),
230 "application/json",
231 )
232 .await
233 }
234
235 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 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 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 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"), ¶ms, None, None).await
312 }
313
314 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 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 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 let backoff = min_delay.saturating_mul(2u32.saturating_pow(attempt));
446 sleep_public(crate::http_client::randomized_delay(backoff)).await;
447 }
448 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 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 ¶ms,
479 Some(body),
480 "application/json",
481 )
482 .await
483 }
484
485 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 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"), ¶ms).await
525 }
526
527 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 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 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 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
613pub struct RequestQueueRequestsIterator {
619 client: RequestQueueClient,
620 page_limit: Option<i64>,
621 buffer: std::collections::VecDeque<RequestQueueRequest>,
622 next_cursor: Option<String>,
624 exhausted: bool,
625}
626
627impl RequestQueueRequestsIterator {
628 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 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 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 #[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 #[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 #[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 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 #[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 #[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 #[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 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}