mod common;
use apify_client::models::RequestQueueRequest;
use serde_json::json;
#[tokio::test(flavor = "multi_thread")]
async fn list_request_queues() {
let client = require_client!();
let page = client
.request_queues()
.list(Default::default())
.await
.expect("listing request queues should succeed");
assert!(page.total >= 0);
}
#[tokio::test(flavor = "multi_thread")]
async fn get_request_queue() {
let client = require_client!();
let name = common::unique_name("rq-get");
let queue = client
.request_queues()
.get_or_create(Some(&name))
.await
.expect("create queue");
let cleanup_client = client.clone();
let id = queue.id.clone();
let _guard = common::Cleanup::new(move || async move {
let _ = cleanup_client.request_queue(&id).delete().await;
});
let fetched = client
.request_queue(&queue.id)
.get()
.await
.expect("get queue by id")
.expect("queue should exist");
assert_eq!(fetched.id, queue.id);
}
#[tokio::test(flavor = "multi_thread")]
async fn iterate_request_queues() {
let client = require_client!();
let name = common::unique_name("rq-iter");
let queue = client
.request_queues()
.get_or_create(Some(&name))
.await
.expect("create queue");
let cleanup_client = client.clone();
let id = queue.id.clone();
let _guard = common::Cleanup::new(move || async move {
let _ = cleanup_client.request_queue(&id).delete().await;
});
let target = queue.id.clone();
assert!(
common::iter_contains_eventually(
|| {
client
.request_queues()
.iterate(apify_client::StorageListOptions {
desc: Some(true),
..Default::default()
})
.with_chunk_size(5)
},
move |q| q.id == target,
)
.await,
"request queue iteration should yield the created queue"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn request_queue_crud_flow() {
let client = require_client!();
let name = common::unique_name("rq");
let queue = client
.request_queues()
.get_or_create(Some(&name))
.await
.expect("create queue");
assert_eq!(queue.name.as_deref(), Some(name.as_str()));
let cleanup_client = client.clone();
let cleanup_id = queue.id.clone();
let _guard = common::Cleanup::new(move || async move {
let _ = cleanup_client.request_queue(&cleanup_id).delete().await;
});
let queue_client = client.request_queue(&queue.id);
assert!(queue_client.get().await.expect("get queue").is_some());
let request = RequestQueueRequest {
id: None,
url: "https://example.com/".to_string(),
unique_key: Some("example".to_string()),
method: Some("GET".to_string()),
user_data: Some(json!({ "label": "START" })),
extra: Default::default(),
};
let added = queue_client
.add_request(&request, false)
.await
.expect("add request");
assert!(!added.request_id.is_empty());
let fetched = queue_client
.get_request(&added.request_id)
.await
.expect("get request")
.expect("request should exist");
assert_eq!(fetched.url, "https://example.com/");
let head = queue_client.list_head(Some(10)).await.expect("list head");
assert!(head
.items
.iter()
.any(|r| r.id.as_deref() == Some(added.request_id.as_str())));
let renamed = common::unique_name("rq-renamed");
let updated = queue_client
.update(&json!({ "name": renamed }))
.await
.expect("update queue");
assert_eq!(updated.name.as_deref(), Some(renamed.as_str()));
queue_client
.delete_request(&added.request_id)
.await
.expect("delete request");
queue_client.delete().await.expect("delete queue");
assert!(queue_client
.get()
.await
.expect("get after delete")
.is_none());
}
#[tokio::test(flavor = "multi_thread")]
async fn request_queue_paginate_multiple_pages() {
let client = require_client!();
let name = common::unique_name("rq-page");
let queue = client
.request_queues()
.get_or_create(Some(&name))
.await
.expect("create queue");
let cleanup_client = client.clone();
let cleanup_id = queue.id.clone();
let _guard = common::Cleanup::new(move || async move {
let _ = cleanup_client.request_queue(&cleanup_id).delete().await;
});
let queue_client = client.request_queue(&queue.id);
const TOTAL: usize = 5;
let mut expected_urls = std::collections::HashSet::new();
for i in 0..TOTAL {
let url = format!("https://example.com/page/{i}");
let request = RequestQueueRequest {
id: None,
url: url.clone(),
unique_key: Some(format!("page-{i}")),
method: Some("GET".to_string()),
user_data: None,
extra: Default::default(),
};
queue_client
.add_request(&request, false)
.await
.expect("add request");
expected_urls.insert(url);
}
let mut iter = queue_client.paginate_requests(Some(2));
let mut seen = std::collections::HashSet::new();
while let Some(req) = iter.next().await.expect("paginate page") {
assert!(
seen.insert(req.url.clone()),
"request {} yielded more than once (pagination cursor bug)",
req.url
);
}
assert_eq!(
seen, expected_urls,
"pagination must yield every added request exactly once across pages"
);
queue_client.delete().await.expect("delete queue");
}
#[tokio::test(flavor = "multi_thread")]
async fn request_queue_batch_add_and_delete() {
let client = require_client!();
let name = common::unique_name("rq-batch");
let queue = client
.request_queues()
.get_or_create(Some(&name))
.await
.expect("create queue");
let cleanup_client = client.clone();
let cleanup_id = queue.id.clone();
let _guard = common::Cleanup::new(move || async move {
let _ = cleanup_client.request_queue(&cleanup_id).delete().await;
});
let queue_client = client.request_queue(&queue.id);
const TOTAL: usize = 8;
let requests: Vec<RequestQueueRequest> = (0..TOTAL)
.map(|i| RequestQueueRequest {
id: None,
url: format!("https://example.com/batch/{i}"),
unique_key: Some(format!("batch-{i}")),
method: Some("GET".to_string()),
user_data: None,
extra: Default::default(),
})
.collect();
let add_result = queue_client
.batch_add_requests(&requests, Default::default())
.await
.expect("batch add requests");
assert_eq!(
add_result.processed_requests.len(),
TOTAL,
"every request should be processed (none rate-limited on a fresh queue)"
);
assert!(add_result.unprocessed_requests.is_empty());
let mut processed_keys: std::collections::HashSet<String> = add_result
.processed_requests
.iter()
.filter_map(|r| r.unique_key.clone())
.collect();
for i in 0..TOTAL {
assert!(
processed_keys.remove(&format!("batch-{i}")),
"processed_requests should report unique_key batch-{i}"
);
}
let delete_by_key: Vec<serde_json::Value> = requests
.iter()
.map(|r| json!({ "uniqueKey": r.unique_key }))
.collect();
let delete_result = queue_client
.batch_delete_requests(&delete_by_key)
.await
.expect("batch delete requests");
assert_eq!(delete_result.processed_requests.len(), TOTAL);
let too_many: Vec<serde_json::Value> = (0..26)
.map(|i| json!({ "uniqueKey": format!("toomany-{i}") }))
.collect();
let err = queue_client
.batch_delete_requests(&too_many)
.await
.expect_err("more than 25 requests must be rejected client-side");
assert!(matches!(
err,
apify_client::ApifyClientError::InvalidArgument(_)
));
queue_client.delete().await.expect("delete queue");
}
#[tokio::test(flavor = "multi_thread")]
async fn request_queue_lock_lifecycle() {
let client = require_client!();
let name = common::unique_name("rq-lock");
let queue = client
.request_queues()
.get_or_create(Some(&name))
.await
.expect("create queue");
let cleanup_client = client.clone();
let cleanup_id = queue.id.clone();
let _guard = common::Cleanup::new(move || async move {
let _ = cleanup_client.request_queue(&cleanup_id).delete().await;
});
let queue_client = client
.request_queue(&queue.id)
.with_client_key("rust-test-client");
let request = RequestQueueRequest {
id: None,
url: "https://example.com/lock".to_string(),
unique_key: Some("lock-example".to_string()),
method: Some("GET".to_string()),
user_data: None,
extra: Default::default(),
};
let added = queue_client
.add_request(&request, false)
.await
.expect("add request");
let listed = queue_client
.list_requests(apify_client::ListRequestsOptions {
limit: Some(10),
..Default::default()
})
.await
.expect("list requests");
assert!(!listed.items.is_empty());
let filtered = queue_client
.list_requests(apify_client::ListRequestsOptions {
limit: Some(10),
filter: Some(vec!["locked".to_string(), "pending".to_string()]),
..Default::default()
})
.await
.expect("list requests with filter");
assert_eq!(
filtered.items.len(),
listed.items.len(),
"filtering on both locked and pending must match the unfiltered listing"
);
let mut iter = queue_client.paginate_requests(Some(10));
let first = iter.next().await.expect("paginate requests");
assert!(first.is_some(), "pagination should yield the added request");
let locked = queue_client
.list_and_lock_head(30, Some(5))
.await
.expect("lock head");
assert!(locked.lock_secs > 0);
queue_client
.prolong_request_lock(&added.request_id, 60, false)
.await
.expect("prolong lock");
queue_client
.delete_request_lock(&added.request_id, false)
.await
.expect("delete lock");
queue_client
.unlock_requests()
.await
.expect("unlock requests");
queue_client.delete().await.expect("delete queue");
}