fix: reject repeated job pagination tokens (#4139)

Fixes #4138

`RemoteDatabase::list_jobs` followed every returned pagination token
without remembering previously seen values. A server-side token cycle
therefore caused repeated requests and duplicate accumulation until the
100-page safeguard returned partial results as a success.

This change tracks non-empty job-list page tokens and returns an
HTTP-context error as soon as a token repeats, matching the existing
`list_functions` behavior. A mock-handler regression test verifies a
repeated `loop` token is rejected after two requests.

Validation:
- `cargo test --quiet --features remote -p lancedb test_list_jobs`
- `cargo fmt --all`
- `cargo check --quiet --features remote --tests --examples`

<!-- lance-gatekeeper-fix:v1 agent=ddd56389737a405d32cf4f43c7696032
generation=1 -->

Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com>
Co-authored-by: Xuanwo <github@xuanwo.io>
This commit is contained in:
lancedb-gatefixer[bot]
2026-09-14 20:50:54 +08:00
committed by GitHub
co-authored by Xuanwo
parent 255da8a854
commit 6bb64c3edb
+43 -3
View File
@@ -700,6 +700,7 @@ impl<S: HttpSend> Database for RemoteDatabase<S> {
async fn list_jobs(&self) -> Result<Vec<JobInfo>> {
let mut out = Vec::new();
let mut page_token: Option<String> = None;
let mut seen_page_tokens = HashSet::new();
for page in 0..MAX_LIST_JOBS_PAGES {
let mut body = serde_json::json!({});
if let Some(token) = &page_token {
@@ -708,7 +709,8 @@ impl<S: HttpSend> Database for RemoteDatabase<S> {
let req = self.client.post("/v1/jobs/list").json(&body);
let (request_id, rsp) = self.client.send(req).await?;
let rsp = self.client.check_response(&request_id, rsp).await?;
let body: RemoteListJobsResponse = rsp.json().await.err_to_http(request_id)?;
let status = rsp.status();
let body: RemoteListJobsResponse = rsp.json().await.err_to_http(request_id.clone())?;
out.extend(body.jobs.into_iter().map(|row| JobInfo {
job_id: row.job_id,
table: row.table,
@@ -716,10 +718,17 @@ impl<S: HttpSend> Database for RemoteDatabase<S> {
state: job_state_to_client(&row.state),
created_at_millis: row.created_at_millis,
}));
page_token = body.page_token;
if page_token.is_none() {
let Some(next_page_token) = body.page_token.filter(|token| !token.is_empty()) else {
break;
};
if !seen_page_tokens.insert(next_page_token.clone()) {
return Err(Error::Http {
source: "Job listing response repeated a page_token".into(),
request_id,
status_code: Some(status),
});
}
page_token = Some(next_page_token);
if page + 1 == MAX_LIST_JOBS_PAGES {
log::warn!(
"list_jobs truncated after {} pages ({} jobs)",
@@ -2634,6 +2643,37 @@ mod tests {
assert_eq!(jobs[2].state, "failed");
}
#[tokio::test]
async fn test_list_jobs_rejects_a_page_token_cycle() {
let requests = Arc::new(AtomicUsize::new(0));
let seen = requests.clone();
let conn = Connection::new_with_handler(move |request| {
let body: serde_json::Value =
serde_json::from_slice(request.body().unwrap().as_bytes().unwrap()).unwrap();
match seen.fetch_add(1, Ordering::SeqCst) {
0 => assert!(body.get("page_token").is_none()),
_ => assert_eq!(body["page_token"], "loop"),
}
http::Response::builder()
.status(200)
.body(r#"{"jobs": [], "page_token": "loop"}"#)
.unwrap()
});
let error = conn.list_jobs().await.unwrap_err();
assert!(
matches!(
&error,
Error::Http {
status_code: Some(http::StatusCode::OK),
..
}
),
"got {error:?}"
);
assert_eq!(requests.load(Ordering::SeqCst), 2);
}
#[tokio::test]
async fn test_open_job() {
let conn = Connection::new_with_handler(|request| {