use serde::de::DeserializeOwned;
use serde::Deserialize;
use tracing::{debug, warn};
use crate::error::Result;
use crate::ApiClient;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Continuation {
Url(String),
Token { param: String, value: String },
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct PageInfo {
pub total: Option<u64>,
pub truncated: bool,
pub next: Option<Continuation>,
}
pub trait Page: DeserializeOwned {
type Item;
fn into_parts(self) -> (Vec<Self::Item>, Option<Continuation>, Option<u64>);
}
#[derive(Debug, Clone, Deserialize)]
pub struct BitbucketPage<T> {
pub values: Vec<T>,
#[serde(default)]
pub next: Option<String>,
#[serde(default)]
pub size: Option<u64>,
}
impl<T: DeserializeOwned> Page for BitbucketPage<T> {
type Item = T;
fn into_parts(self) -> (Vec<T>, Option<Continuation>, Option<u64>) {
let cursor = self
.next
.filter(|url| !url.trim().is_empty())
.map(Continuation::Url);
(self.values, cursor, self.size)
}
}
#[derive(Debug, Clone, Deserialize)]
pub struct JiraPage<T> {
pub issues: Vec<T>,
#[serde(default, rename = "nextPageToken")]
pub next_page_token: Option<String>,
#[serde(default, rename = "isLast")]
pub is_last: Option<bool>,
}
pub const JIRA_PAGE_TOKEN_PARAM: &str = "nextPageToken";
impl<T: DeserializeOwned> Page for JiraPage<T> {
type Item = T;
fn into_parts(self) -> (Vec<T>, Option<Continuation>, Option<u64>) {
let finished = self.is_last.unwrap_or(false);
let cursor = if finished {
None
} else {
self.next_page_token
.filter(|token| !token.trim().is_empty())
.map(|value| Continuation::Token {
param: JIRA_PAGE_TOKEN_PARAM.to_string(),
value,
})
};
(self.issues, cursor, None)
}
}
#[derive(Debug, Clone, Deserialize)]
pub struct JiraOffsetPage<T> {
pub values: Vec<T>,
#[serde(default, rename = "startAt")]
pub start_at: Option<u64>,
#[serde(default, rename = "maxResults")]
pub max_results: Option<u64>,
#[serde(default)]
pub total: Option<u64>,
#[serde(default, rename = "isLast")]
pub is_last: Option<bool>,
}
pub const JIRA_START_AT_PARAM: &str = "startAt";
impl<T: DeserializeOwned> Page for JiraOffsetPage<T> {
type Item = T;
fn into_parts(self) -> (Vec<T>, Option<Continuation>, Option<u64>) {
let start = self.start_at.unwrap_or(0);
let returned = self.values.len() as u64;
let next_offset = start + returned;
let finished = match self.is_last {
Some(is_last) => is_last,
None => match self.total {
Some(total) => next_offset >= total,
None => returned == 0,
},
};
let cursor = if finished || returned == 0 {
None
} else {
Some(Continuation::Token {
param: JIRA_START_AT_PARAM.to_string(),
value: next_offset.to_string(),
})
};
(self.values, cursor, self.total)
}
}
#[derive(Debug, Clone, Copy)]
pub struct PageLimits {
pub limit: Option<usize>,
pub budget: usize,
}
impl PageLimits {
pub const DEFAULT_BUDGET: usize = 50;
pub fn new(limit: Option<usize>) -> Self {
Self {
limit,
budget: Self::DEFAULT_BUDGET,
}
}
pub fn from_cli_limit(limit: usize) -> Self {
Self::new(if limit == 0 { None } else { Some(limit) })
}
pub fn with_budget(mut self, budget: usize) -> Self {
self.budget = budget;
self
}
}
pub async fn fetch_paged<P: Page>(
client: &ApiClient,
path: &str,
limits: PageLimits,
) -> Result<(Vec<P::Item>, PageInfo)> {
let mut items: Vec<P::Item> = Vec::new();
let mut info = PageInfo::default();
let mut request = path.to_string();
for attempt in 0..limits.budget.max(1) {
debug!(request = %request, attempt, "Fetching page");
let page: P = client.get(&request).await?;
let (page_items, cursor, total) = page.into_parts();
if total.is_some() {
info.total = total;
}
let empty_page = page_items.is_empty();
items.extend(page_items);
if let Some(limit) = limits.limit {
if items.len() >= limit {
let dropped = items.len() > limit;
items.truncate(limit);
info.truncated = dropped || cursor.is_some();
info.next = cursor;
return Ok((items, info));
}
}
let Some(cursor) = cursor else {
return Ok((items, info));
};
if empty_page {
warn!("Stopping pagination: the server returned an empty page with a cursor");
info.truncated = true;
info.next = Some(cursor);
return Ok((items, info));
}
if attempt + 1 >= limits.budget.max(1) {
warn!(
budget = limits.budget,
collected = items.len(),
"Stopping pagination: request budget exhausted; the result is incomplete"
);
info.truncated = true;
info.next = Some(cursor);
return Ok((items, info));
}
request = match cursor {
Continuation::Url(url) => url,
Continuation::Token { param, value } => append_query(path, ¶m, &value),
};
}
Ok((items, info))
}
fn append_query(path: &str, key: &str, value: &str) -> String {
let (base, query) = match path.split_once('?') {
Some((base, query)) => (base, Some(query)),
None => (path, None),
};
let mut pairs: Vec<String> = query
.map(|q| {
q.split('&')
.filter(|pair| !pair.is_empty())
.filter(|pair| {
let name = pair.split('=').next().unwrap_or("");
name != key
})
.map(|pair| pair.to_string())
.collect()
})
.unwrap_or_default();
pairs.push(format!(
"{}={}",
encode_query_component(key),
encode_query_component(value)
));
format!("{base}?{}", pairs.join("&"))
}
fn encode_query_component(value: &str) -> String {
let mut out = String::with_capacity(value.len());
for byte in value.as_bytes() {
match byte {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'.' | b'_' | b'~' => {
out.push(*byte as char)
}
_ => out.push_str(&format!("%{byte:02X}")),
}
}
out
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[derive(Debug, Deserialize, PartialEq)]
struct Item {
name: String,
}
fn bb_page(value: serde_json::Value) -> BitbucketPage<Item> {
serde_json::from_value(value).unwrap()
}
fn jira_page(value: serde_json::Value) -> JiraPage<Item> {
serde_json::from_value(value).unwrap()
}
#[test]
fn a_bitbucket_page_yields_its_next_url() {
let page = bb_page(json!({
"values": [{"name": "a"}],
"next": "https://api.bitbucket.org/2.0/x?page=2",
"size": 7
}));
let (items, cursor, total) = page.into_parts();
assert_eq!(items.len(), 1);
assert_eq!(
cursor,
Some(Continuation::Url(
"https://api.bitbucket.org/2.0/x?page=2".to_string()
))
);
assert_eq!(total, Some(7));
}
#[test]
fn a_bitbucket_page_without_next_is_the_last() {
let page = bb_page(json!({"values": [{"name": "a"}]}));
assert_eq!(page.into_parts().1, None);
}
#[test]
fn a_blank_next_is_not_a_cursor() {
let page = bb_page(json!({"values": [], "next": " "}));
assert_eq!(page.into_parts().1, None);
}
#[test]
fn is_last_overrides_a_trailing_jira_token() {
let page = jira_page(json!({
"issues": [{"name": "a"}],
"nextPageToken": "abc",
"isLast": true
}));
assert_eq!(page.into_parts().1, None);
}
#[test]
fn a_jira_token_carries_its_parameter_name() {
let page = jira_page(json!({"issues": [], "nextPageToken": "abc"}));
assert_eq!(
page.into_parts().1,
Some(Continuation::Token {
param: JIRA_PAGE_TOKEN_PARAM.to_string(),
value: "abc".to_string()
})
);
}
#[test]
fn jira_reports_no_total() {
let page = jira_page(json!({"issues": [{"name": "a"}]}));
assert_eq!(page.into_parts().2, None);
}
use wiremock::matchers::{method, path as path_matcher, query_param};
use wiremock::{Mock, MockServer, ResponseTemplate};
fn client_for(server: &MockServer) -> ApiClient {
ApiClient::new(server.uri()).unwrap()
}
#[tokio::test]
async fn a_bitbucket_collection_is_followed_to_the_end() {
let server = MockServer::start().await;
let page_two = format!("{}/items?page=2", server.uri());
Mock::given(method("GET"))
.and(path_matcher("/items"))
.and(query_param("page", "2"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"values": [{"name": "c"}]
})))
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path_matcher("/items"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"values": [{"name": "a"}, {"name": "b"}],
"next": page_two,
"size": 3
})))
.mount(&server)
.await;
let (items, info) = fetch_paged::<BitbucketPage<Item>>(
&client_for(&server),
"/items",
PageLimits::new(None),
)
.await
.unwrap();
assert_eq!(items.len(), 3, "every page must be collected");
assert_eq!(items[2].name, "c");
assert!(!info.truncated, "a complete result is not truncated");
assert_eq!(info.total, Some(3));
}
#[tokio::test]
async fn a_limit_truncates_and_says_so() {
let server = MockServer::start().await;
let page_two = format!("{}/items?page=2", server.uri());
Mock::given(method("GET"))
.and(path_matcher("/items"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"values": [{"name": "a"}, {"name": "b"}],
"next": page_two
})))
.mount(&server)
.await;
let (items, info) = fetch_paged::<BitbucketPage<Item>>(
&client_for(&server),
"/items",
PageLimits::new(Some(1)),
)
.await
.unwrap();
assert_eq!(items.len(), 1);
assert!(info.truncated, "a capped result must be labelled truncated");
assert!(info.next.is_some(), "and must say where it stopped");
}
#[tokio::test]
async fn the_budget_bounds_a_server_that_never_ends() {
let server = MockServer::start().await;
let forever = format!("{}/items?page=next", server.uri());
Mock::given(method("GET"))
.and(path_matcher("/items"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"values": [{"name": "x"}],
"next": forever
})))
.mount(&server)
.await;
let (items, info) = fetch_paged::<BitbucketPage<Item>>(
&client_for(&server),
"/items",
PageLimits::new(None).with_budget(3),
)
.await
.unwrap();
assert_eq!(items.len(), 3, "one item per allowed request");
assert!(
info.truncated,
"budget exhaustion is truncation, not success"
);
}
#[tokio::test]
async fn an_empty_page_with_a_cursor_stops_the_walk() {
let server = MockServer::start().await;
let forever = format!("{}/items?page=next", server.uri());
Mock::given(method("GET"))
.and(path_matcher("/items"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"values": [],
"next": forever
})))
.mount(&server)
.await;
let (items, _) = fetch_paged::<BitbucketPage<Item>>(
&client_for(&server),
"/items",
PageLimits::new(None).with_budget(20),
)
.await
.unwrap();
assert!(items.is_empty());
let requests = server.received_requests().await.unwrap_or_default();
assert_eq!(
requests.len(),
1,
"must not keep asking: {}",
requests.len()
);
}
#[tokio::test]
async fn an_offset_paged_endpoint_is_followed_to_the_end() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_matcher("/project/search"))
.and(query_param("startAt", "2"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"values": [{"name": "c"}], "startAt": 2, "total": 3
})))
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path_matcher("/project/search"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"values": [{"name": "a"}, {"name": "b"}], "startAt": 0, "total": 3
})))
.mount(&server)
.await;
let (items, info) = fetch_paged::<JiraOffsetPage<Item>>(
&client_for(&server),
"/project/search?expand=lead",
PageLimits::new(None),
)
.await
.unwrap();
assert_eq!(items.len(), 3, "every page collected");
assert!(!info.truncated);
assert_eq!(info.total, Some(3));
for request in server.received_requests().await.unwrap_or_default() {
let query = request.url.query().unwrap_or("");
assert!(
query.matches("startAt").count() <= 1,
"offset accumulated: {query}"
);
assert!(
query.contains("expand=lead"),
"the original query must survive: {query}"
);
}
}
#[tokio::test]
async fn a_jira_token_never_accumulates() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_matcher("/search"))
.and(query_param("nextPageToken", "t2"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"issues": [{"name": "c"}], "isLast": true
})))
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path_matcher("/search"))
.and(query_param("nextPageToken", "t1"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"issues": [{"name": "b"}], "nextPageToken": "t2", "isLast": false
})))
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path_matcher("/search"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"issues": [{"name": "a"}], "nextPageToken": "t1", "isLast": false
})))
.mount(&server)
.await;
let (items, info) = fetch_paged::<JiraPage<Item>>(
&client_for(&server),
"/search?jql=project%3DX",
PageLimits::new(None),
)
.await
.unwrap();
assert_eq!(items.len(), 3, "all three pages");
assert!(!info.truncated);
for request in server.received_requests().await.unwrap_or_default() {
let query = request.url.query().unwrap_or("");
assert!(
query.matches("nextPageToken").count() <= 1,
"token accumulated: {query}"
);
assert!(
query.contains("jql=project"),
"the original query must survive: {query}"
);
}
}
#[tokio::test]
async fn a_foreign_next_url_is_refused() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_matcher("/items"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"values": [{"name": "a"}],
"next": "https://evil.example.com/2.0/items?page=2"
})))
.mount(&server)
.await;
let result = fetch_paged::<BitbucketPage<Item>>(
&client_for(&server),
"/items",
PageLimits::new(None),
)
.await;
assert!(
result.is_err(),
"a cross-origin cursor must not be followed"
);
}
#[tokio::test]
async fn overshooting_the_limit_on_a_final_page_is_still_truncation() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_matcher("/search"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"issues": [{"name": "a"}, {"name": "b"}, {"name": "c"}],
"isLast": true
})))
.mount(&server)
.await;
let (items, info) = fetch_paged::<JiraPage<Item>>(
&client_for(&server),
"/search",
PageLimits::new(Some(2)),
)
.await
.unwrap();
assert_eq!(items.len(), 2);
assert!(
info.truncated,
"a page that overshot the limit dropped rows and must say so"
);
}
#[tokio::test]
async fn hitting_the_limit_exactly_is_not_truncation() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_matcher("/search"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"issues": [{"name": "a"}, {"name": "b"}],
"isLast": true
})))
.mount(&server)
.await;
let (items, info) = fetch_paged::<JiraPage<Item>>(
&client_for(&server),
"/search",
PageLimits::new(Some(2)),
)
.await
.unwrap();
assert_eq!(items.len(), 2);
assert!(!info.truncated, "nothing was dropped: {info:?}");
}
#[tokio::test]
async fn an_abandoned_walk_is_reported_as_truncated() {
let server = MockServer::start().await;
let forever = format!("{}/items?page=next", server.uri());
Mock::given(method("GET"))
.and(path_matcher("/items"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"values": [],
"next": forever
})))
.mount(&server)
.await;
let (_, info) = fetch_paged::<BitbucketPage<Item>>(
&client_for(&server),
"/items",
PageLimits::new(None),
)
.await
.unwrap();
assert!(info.truncated, "the server said more exists: {info:?}");
}
#[tokio::test]
async fn a_body_without_the_items_key_is_an_error() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path_matcher("/items"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({"page": 1})))
.mount(&server)
.await;
let result = fetch_paged::<BitbucketPage<Item>>(
&client_for(&server),
"/items",
PageLimits::new(None),
)
.await;
assert!(
result.is_err(),
"a missing values key must not read as an empty result"
);
}
#[test]
fn a_zero_cli_limit_means_no_limit() {
assert_eq!(PageLimits::from_cli_limit(0).limit, None);
assert_eq!(PageLimits::from_cli_limit(25).limit, Some(25));
}
fn offset_page(value: serde_json::Value) -> JiraOffsetPage<Item> {
serde_json::from_value(value).unwrap()
}
#[test]
fn an_offset_page_advances_by_what_it_returned() {
let page = offset_page(json!({
"values": [{"name": "a"}, {"name": "b"}],
"startAt": 0, "maxResults": 2, "total": 5
}));
let (items, cursor, total) = page.into_parts();
assert_eq!(items.len(), 2);
assert_eq!(
cursor,
Some(Continuation::Token {
param: JIRA_START_AT_PARAM.to_string(),
value: "2".to_string()
})
);
assert_eq!(total, Some(5));
}
#[test]
fn is_last_ends_an_offset_walk() {
let page = offset_page(json!({
"values": [{"name": "a"}], "startAt": 0, "total": 99, "isLast": true
}));
assert_eq!(page.into_parts().1, None);
}
#[test]
fn reaching_the_total_ends_an_offset_walk() {
let page = offset_page(json!({
"values": [{"name": "a"}], "startAt": 4, "total": 5
}));
assert_eq!(page.into_parts().1, None);
}
#[test]
fn an_empty_offset_page_ends_the_walk() {
let page = offset_page(json!({"values": [], "startAt": 10}));
assert_eq!(page.into_parts().1, None);
}
#[test]
fn append_query_adds_a_parameter() {
assert_eq!(append_query("/x", "t", "1"), "/x?t=1");
assert_eq!(append_query("/x?a=b", "t", "1"), "/x?a=b&t=1");
}
#[test]
fn append_query_replaces_rather_than_duplicating() {
let once = append_query("/x?a=b", "t", "1");
let twice = append_query(&once, "t", "2");
assert_eq!(twice, "/x?a=b&t=2");
assert_eq!(twice.matches("t=").count(), 1);
}
#[test]
fn append_query_encodes_its_value() {
assert_eq!(append_query("/x", "t", "a b&c"), "/x?t=a%20b%26c");
}
}