use std::collections::HashSet;
use chrono::{DateTime, Timelike, Utc};
use tracing::warn;
pub type PagedItem = (String, Option<DateTime<Utc>>);
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PageRequest {
pub since: Option<DateTime<Utc>>,
pub next_page_token: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PageStep {
pub is_new: Vec<bool>,
pub more: bool,
}
#[derive(Debug)]
pub struct KeysetPager {
since: Option<DateTime<Utc>>,
next_page_token: Option<String>,
seen: HashSet<String>,
pages: usize,
max_pages: usize,
offset_paged_minute: Option<DateTime<Utc>>,
}
impl KeysetPager {
pub fn new(since: Option<DateTime<Utc>>, max_pages: usize) -> Self {
Self {
since,
next_page_token: None,
seen: HashSet::new(),
pages: 0,
max_pages: max_pages.max(1),
offset_paged_minute: None,
}
}
pub fn request(&self) -> PageRequest {
PageRequest {
since: self.since,
next_page_token: self.next_page_token.clone(),
}
}
pub fn offset_paged_minute(&self) -> Option<DateTime<Utc>> {
self.offset_paged_minute
}
pub fn record_page(&mut self, items: &[PagedItem], next_page_token: Option<&str>) -> PageStep {
self.pages += 1;
let mut is_new = Vec::with_capacity(items.len());
let mut page_max: Option<DateTime<Utc>> = None;
for (key, updated) in items {
is_new.push(self.seen.insert(key.clone()));
if let Some(u) = updated {
page_max = Some(page_max.map_or(*u, |m: DateTime<Utc>| m.max(*u)));
}
}
let Some(token) = next_page_token else {
return PageStep {
is_new,
more: false,
};
};
if self.pages >= self.max_pages {
warn!(
pages = self.pages,
"JIRA changelog walk hit its page cap; truncating this run \
(the sync cursor only advances over ingested tickets, so the \
next run resumes here)"
);
return PageStep {
is_new,
more: false,
};
}
match page_max {
Some(max) if minute(Some(max)) > minute(self.since) => {
self.since = Some(max);
self.next_page_token = None;
}
_ => {
let stalled = minute(self.since).or_else(|| minute(page_max));
self.offset_paged_minute = match (self.offset_paged_minute, stalled) {
(Some(existing), Some(new)) => Some(existing.min(new)),
(existing, new) => existing.or(new),
};
self.next_page_token = Some(token.to_string());
}
}
PageStep { is_new, more: true }
}
}
fn minute(dt: Option<DateTime<Utc>>) -> Option<DateTime<Utc>> {
dt.map(|d| {
d.with_second(0)
.and_then(|d| d.with_nanosecond(0))
.unwrap_or(d)
})
}
#[cfg(test)]
mod tests {
use super::*;
fn dt(s: &str) -> DateTime<Utc> {
DateTime::parse_from_rfc3339(s)
.expect("valid rfc3339 fixture")
.with_timezone(&Utc)
}
fn item(key: &str, updated: &str) -> PagedItem {
(key.to_string(), Some(dt(updated)))
}
#[test]
fn pager_starts_at_the_scope_lower_bound() {
let pager = KeysetPager::new(Some(dt("2026-01-01T00:00:00Z")), 10);
assert_eq!(
pager.request(),
PageRequest {
since: Some(dt("2026-01-01T00:00:00Z")),
next_page_token: None,
}
);
}
#[test]
fn pager_stops_when_the_server_returns_no_token() {
let mut pager = KeysetPager::new(None, 10);
let step = pager.record_page(&[item("P-1", "2026-01-01T00:00:00Z")], None);
assert!(!step.more, "an absent continuation token ends the walk");
assert_eq!(step.is_new, vec![true]);
}
#[test]
fn pager_does_not_stop_on_a_short_page() {
let mut pager = KeysetPager::new(None, 10);
let step = pager.record_page(&[item("P-1", "2026-01-01T00:00:00Z")], Some("t1"));
assert!(
step.more,
"a short page carrying a token must not end the walk"
);
}
#[test]
fn pager_reanchors_window_on_a_later_minute() {
let mut pager = KeysetPager::new(Some(dt("2026-01-01T00:00:00Z")), 10);
let page = vec![
item("P-1", "2026-01-01T00:01:00Z"),
item("P-2", "2026-01-01T00:02:00Z"),
];
let step = pager.record_page(&page, Some("t1"));
assert!(step.more);
assert_eq!(
pager.request(),
PageRequest {
since: Some(dt("2026-01-01T00:02:00Z")),
next_page_token: None,
},
"the next window must start at the page's max updated, not at a \
position carried over from a different query"
);
}
#[test]
fn pager_falls_back_to_the_page_cursor_within_one_minute() {
let mut pager = KeysetPager::new(Some(dt("2026-01-01T00:00:00Z")), 10);
let page = vec![
item("P-1", "2026-01-01T00:00:10Z"),
item("P-2", "2026-01-01T00:00:20Z"),
];
let step = pager.record_page(&page, Some("t1"));
assert!(step.more);
assert_eq!(
pager.request(),
PageRequest {
since: Some(dt("2026-01-01T00:00:00Z")),
next_page_token: Some("t1".to_string()),
},
"a page confined to the window's own minute must advance by the \
server's continuation token"
);
}
#[test]
fn pager_deduplicates_reread_boundary_items() {
let mut pager = KeysetPager::new(None, 10);
let first = vec![
item("P-1", "2026-01-01T00:01:00Z"),
item("P-2", "2026-01-01T00:02:00Z"),
];
pager.record_page(&first, Some("t1"));
let second = vec![
item("P-2", "2026-01-01T00:02:00Z"),
item("P-3", "2026-01-01T00:03:00Z"),
];
let step = pager.record_page(&second, Some("t2"));
assert_eq!(step.is_new, vec![false, true]);
}
#[test]
fn pager_stops_at_max_pages() {
let mut pager = KeysetPager::new(None, 2);
let page = vec![
item("P-1", "2026-01-01T00:01:00Z"),
item("P-2", "2026-01-01T00:02:00Z"),
];
assert!(pager.record_page(&page, Some("t1")).more);
let page2 = vec![
item("P-3", "2026-01-01T00:03:00Z"),
item("P-4", "2026-01-01T00:04:00Z"),
];
assert!(
!pager.record_page(&page2, Some("t2")).more,
"the page cap must stop the walk rather than loop forever"
);
}
#[test]
fn pager_reports_the_minute_it_offset_paged() {
let mut pager = KeysetPager::new(Some(dt("2026-01-01T00:00:30Z")), 10);
assert_eq!(pager.offset_paged_minute(), None);
let page = vec![
item("P-1", "2026-01-01T00:00:40Z"),
item("P-2", "2026-01-01T00:00:50Z"),
];
pager.record_page(&page, Some("t1"));
assert_eq!(
pager.offset_paged_minute(),
Some(dt("2026-01-01T00:00:00Z")),
"the reported value is the start of the stalled minute"
);
}
#[test]
fn pager_reports_no_offset_minute_on_a_clean_walk() {
let mut pager = KeysetPager::new(None, 10);
let page = vec![
item("P-1", "2026-01-01T00:01:00Z"),
item("P-2", "2026-01-01T00:02:00Z"),
];
pager.record_page(&page, Some("t1"));
assert_eq!(pager.offset_paged_minute(), None);
}
#[test]
fn pager_advances_by_cursor_when_a_page_has_no_timestamps() {
let mut pager = KeysetPager::new(Some(dt("2026-01-01T00:00:00Z")), 10);
let page = vec![("P-1".to_string(), None), ("P-2".to_string(), None)];
let step = pager.record_page(&page, Some("t1"));
assert!(step.more);
assert_eq!(pager.request().next_page_token.as_deref(), Some("t1"));
}
}