diff --git a/CHANGELOG.md b/CHANGELOG.md index 2ce86054..9d111d9c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,6 +14,7 @@ If you are upgrading from v0.16.x, replace the binary (or run `docker pull`). If - Live tracing in community and OSS versions. - Timezone changes from the `AccountSettings` object return `invalidProperties`. - `mail-parser` panic with certain messages containing corrupted attachments. +- Pagination by anchor for queued messages, tasks and metrics. ## [0.16.3] - 2026-04-30 diff --git a/crates/jmap/src/registry/mapping/queued_message.rs b/crates/jmap/src/registry/mapping/queued_message.rs index 1f3501c5..40dc233a 100644 --- a/crates/jmap/src/registry/mapping/queued_message.rs +++ b/crates/jmap/src/registry/mapping/queued_message.rs @@ -472,9 +472,18 @@ pub(crate) async fn queued_message_query( let mut total = 0; if let Some(anchor) = req.request.anchor { - let anchor = anchor.id(); - if anchor > due_from { - due_from = anchor; + let anchor_id = anchor.id(); + if let Some(archive) = req.server.read_message_archive(anchor_id).await? + && let Ok(archived) = archive.unarchive::() + && let Some(anchor_due) = archived.next_delivery_event(queue_name) + && anchor_due >= due_from + && anchor_due <= due_to + { + if params.sort_ascending { + due_from = anchor_due; + } else { + due_to = anchor_due; + } } } diff --git a/crates/jmap/src/registry/mapping/task.rs b/crates/jmap/src/registry/mapping/task.rs index 0d0b97c3..22b390d0 100644 --- a/crates/jmap/src/registry/mapping/task.rs +++ b/crates/jmap/src/registry/mapping/task.rs @@ -399,7 +399,7 @@ pub(crate) async fn task_get( pub(crate) async fn task_query( mut req: RegistryQueryResponse<'_>, ) -> trc::Result { - let mut due_from = 100u64; + let mut due_from = 1u64; let mut due_to = u64::MAX; let mut typ = None; @@ -466,6 +466,7 @@ pub(crate) async fn task_query( .extract_parameters(req.server.core.jmap.query_max_results, None)?; let mut from_id = 0u64; + let mut to_id = u64::MAX; if let Some(anchor_id) = anchor_id && let Some(anchor_task) = req .server @@ -478,8 +479,13 @@ pub(crate) async fn task_query( { let anchor_due = anchor_task.due_timestamp(); if anchor_due >= due_from && anchor_due <= due_to { - due_from = anchor_due; - from_id = anchor_id; + if params.sort_ascending { + due_from = anchor_due; + from_id = anchor_id; + } else { + due_to = anchor_due; + to_id = anchor_id; + } } } @@ -497,7 +503,7 @@ pub(crate) async fn task_query( due: due_from, })); let to_key = ValueKey::from(ValueClass::TaskQueue(TaskQueueClass::Due { - id: u64::MAX, + id: to_id, due: due_to, })); diff --git a/crates/jmap/src/registry/mapping/telemetry.rs b/crates/jmap/src/registry/mapping/telemetry.rs index b0ac3dc7..10091230 100644 --- a/crates/jmap/src/registry/mapping/telemetry.rs +++ b/crates/jmap/src/registry/mapping/telemetry.rs @@ -411,17 +411,21 @@ pub(crate) async fn metric_query( if ts_from != 0 { ts_from = SnowflakeIdGenerator::from_timestamp(ts_from).unwrap_or(0); } - if let Some(anchor) = req.request.anchor { - let anchor = anchor.id(); - if anchor > ts_from { - ts_from = anchor; - } - } - if ts_to != u64::MAX { ts_to = SnowflakeIdGenerator::from_timestamp(ts_to).unwrap_or(u64::MAX); } + if let Some(anchor) = req.request.anchor { + let anchor = anchor.id(); + if params.sort_ascending { + if anchor > ts_from { + ts_from = anchor; + } + } else if anchor < ts_to { + ts_to = anchor; + } + } + let from_key = ValueKey::from(ValueClass::Telemetry(TelemetryClass::Metric(ts_from))); let to_key = ValueKey::from(ValueClass::Telemetry(TelemetryClass::Metric(ts_to))); diff --git a/crates/main/src/test_data.rs b/crates/main/src/test_data.rs index 9d2adb5f..056668cc 100644 --- a/crates/main/src/test_data.rs +++ b/crates/main/src/test_data.rs @@ -64,6 +64,7 @@ pub async fn insert_test_data(server: &Server) { let object_id = ObjectType::TlsInternalReport.to_id(); let item_id = server.inner.data.queue_id_gen.generate(); let mut batch = BatchBuilder::new(); + batch.clear(report.primary_key()); report.write_ops(&mut batch, item_id, true); let report_bytes = report.to_pickled_vec(); batch.set( @@ -77,6 +78,7 @@ pub async fn insert_test_data(server: &Server) { let object_id = ObjectType::DmarcInternalReport.to_id(); let item_id = server.inner.data.queue_id_gen.generate(); let mut batch = BatchBuilder::new(); + batch.clear(report.primary_key()); report.write_ops(&mut batch, item_id, true); let report_bytes = report.to_pickled_vec(); batch.set( diff --git a/tests/src/smtp/management/queue.rs b/tests/src/smtp/management/queue.rs index 4488f109..fa4d1d3a 100644 --- a/tests/src/smtp/management/queue.rs +++ b/tests/src/smtp/management/queue.rs @@ -25,6 +25,7 @@ use registry::{ }; use serde_json::json; use std::time::{Duration, Instant}; +use types::id::Id; #[tokio::test] #[serial_test::serial] @@ -335,6 +336,110 @@ async fn manage_queue() { ); } + // Test pagination (forward and reverse) + let asc_order: Vec = admin + .registry_query_paginated( + ObjectType::QueuedMessage, + "due", + true, + None, + None, + None, + None, + false, + ) + .await + .object_ids() + .collect(); + assert_eq!(asc_order.len(), 6, "expected 6 messages, got {asc_order:?}"); + let desc_order: Vec = asc_order.iter().rev().copied().collect(); + + for chunk_start in [0usize, 2, 4] { + let asc = admin + .registry_query_paginated( + ObjectType::QueuedMessage, + "due", + true, + Some(chunk_start as i32), + Some(2), + None, + None, + false, + ) + .await + .object_ids() + .collect::>(); + assert_eq!( + asc, + asc_order[chunk_start..chunk_start + 2], + "ascending position={chunk_start} limit=2", + ); + + let desc = admin + .registry_query_paginated( + ObjectType::QueuedMessage, + "due", + false, + Some(chunk_start as i32), + Some(2), + None, + None, + false, + ) + .await + .object_ids() + .collect::>(); + assert_eq!( + desc, + desc_order[chunk_start..chunk_start + 2], + "descending position={chunk_start} limit=2", + ); + } + + for anchor_idx in [1usize, 3] { + let asc = admin + .registry_query_paginated( + ObjectType::QueuedMessage, + "due", + true, + None, + Some(2), + Some(asc_order[anchor_idx]), + Some(1), + false, + ) + .await + .object_ids() + .collect::>(); + assert_eq!( + asc, + asc_order[anchor_idx + 1..anchor_idx + 3], + "ascending anchor={} offset=1 limit=2", + asc_order[anchor_idx], + ); + + let desc = admin + .registry_query_paginated( + ObjectType::QueuedMessage, + "due", + false, + None, + Some(2), + Some(desc_order[anchor_idx]), + Some(1), + false, + ) + .await + .object_ids() + .collect::>(); + assert_eq!( + desc, + desc_order[anchor_idx + 1..anchor_idx + 3], + "descending anchor={} offset=1 limit=2", + desc_order[anchor_idx], + ); + } + // Retry delivery admin .registry_update_object( diff --git a/tests/src/system/task.rs b/tests/src/system/task.rs index 665b45dd..cee9b178 100644 --- a/tests/src/system/task.rs +++ b/tests/src/system/task.rs @@ -137,9 +137,145 @@ pub async fn test(test: &mut TestServer) { .await .assert_destroyed(&[task.id]); + pagination_test(test).await; + test.cleanup().await; } +async fn pagination_test(test: &mut TestServer) { + println!("Running Task pagination tests..."); + let admin = test.account("admin@example.org"); + + admin.assert_no_tasks().await; + + let mut created = Vec::with_capacity(12); + for i in 0..12u64 { + created.push(admin.schedule_test_task(TASK_SUCCESS, 3600 + i).await); + } + + let asc_order: Vec = admin + .registry_query_paginated(ObjectType::Task, "due", true, None, None, None, None, false) + .await + .object_ids() + .collect(); + assert_eq!(asc_order.len(), 12, "expected 12 tasks, got {}", asc_order.len()); + + let desc_order: Vec = asc_order.iter().rev().copied().collect(); + + for chunk_start in [0usize, 5, 10] { + let chunk_size = std::cmp::min(5, 12 - chunk_start); + + let asc = admin + .registry_query_paginated( + ObjectType::Task, + "due", + true, + Some(chunk_start as i32), + Some(5), + None, + None, + false, + ) + .await + .object_ids() + .collect::>(); + assert_eq!( + asc, + asc_order[chunk_start..chunk_start + chunk_size], + "ascending position={chunk_start} limit=5", + ); + + let desc = admin + .registry_query_paginated( + ObjectType::Task, + "due", + false, + Some(chunk_start as i32), + Some(5), + None, + None, + false, + ) + .await + .object_ids() + .collect::>(); + assert_eq!( + desc, + desc_order[chunk_start..chunk_start + chunk_size], + "descending position={chunk_start} limit=5", + ); + } + + for anchor_idx in [4usize, 9] { + let chunk_size = std::cmp::min(5, 12 - anchor_idx - 1); + + let asc = admin + .registry_query_paginated( + ObjectType::Task, + "due", + true, + None, + Some(5), + Some(asc_order[anchor_idx]), + Some(1), + false, + ) + .await + .object_ids() + .collect::>(); + assert_eq!( + asc, + asc_order[anchor_idx + 1..anchor_idx + 1 + chunk_size], + "ascending anchor={} offset=1 limit=5", + asc_order[anchor_idx], + ); + + let desc = admin + .registry_query_paginated( + ObjectType::Task, + "due", + false, + None, + Some(5), + Some(desc_order[anchor_idx]), + Some(1), + false, + ) + .await + .object_ids() + .collect::>(); + assert_eq!( + desc, + desc_order[anchor_idx + 1..anchor_idx + 1 + chunk_size], + "descending anchor={} offset=1 limit=5", + desc_order[anchor_idx], + ); + } + + let response = admin + .registry_query_paginated( + ObjectType::Task, + "due", + true, + Some(0), + Some(5), + None, + None, + true, + ) + .await; + let total = response + .pointer("/methodResponses/0/1/total") + .and_then(|v| v.as_u64()); + assert_eq!(total, Some(12), "expected calculateTotal=12"); + + admin + .registry_destroy(ObjectType::Task, created.clone()) + .await + .assert_destroyed(&created); + admin.assert_no_tasks().await; +} + impl Account { async fn schedule_test_task(&self, test_type: u64, schedule_in: u64) -> Id { self.registry_create_object(Task::StoreMaintenance(TaskStoreMaintenance { diff --git a/tests/src/telemetry/metrics.rs b/tests/src/telemetry/metrics.rs index 800aee4c..d966f4ab 100644 --- a/tests/src/telemetry/metrics.rs +++ b/tests/src/telemetry/metrics.rs @@ -72,6 +72,114 @@ pub async fn test(test: &TestServer) { metric_ids.len() ); + // Test pagination (forward and reverse) + let asc_order: Vec = admin + .registry_query_paginated( + ObjectType::Metric, + "timestamp", + true, + None, + None, + None, + None, + false, + ) + .await + .object_ids() + .collect(); + assert!(asc_order.len() > 100, "expected >100 metrics, got {}", asc_order.len()); + let desc_order: Vec = asc_order.iter().rev().copied().collect(); + let total = asc_order.len(); + let limit = 25usize; + + for chunk_start in [0usize, limit, total - limit] { + let asc = admin + .registry_query_paginated( + ObjectType::Metric, + "timestamp", + true, + Some(chunk_start as i32), + Some(limit), + None, + None, + false, + ) + .await + .object_ids() + .collect::>(); + assert_eq!( + asc, + asc_order[chunk_start..chunk_start + limit], + "ascending position={chunk_start} limit={limit}", + ); + + let desc = admin + .registry_query_paginated( + ObjectType::Metric, + "timestamp", + false, + Some(chunk_start as i32), + Some(limit), + None, + None, + false, + ) + .await + .object_ids() + .collect::>(); + assert_eq!( + desc, + desc_order[chunk_start..chunk_start + limit], + "descending position={chunk_start} limit={limit}", + ); + } + + for anchor_idx in [limit - 1, total - limit - 1] { + let asc = admin + .registry_query_paginated( + ObjectType::Metric, + "timestamp", + true, + None, + Some(limit), + Some(asc_order[anchor_idx]), + Some(1), + false, + ) + .await + .object_ids() + .collect::>(); + let asc_size = std::cmp::min(limit, total - anchor_idx - 1); + assert_eq!( + asc, + asc_order[anchor_idx + 1..anchor_idx + 1 + asc_size], + "ascending anchor={} offset=1 limit={limit}", + asc_order[anchor_idx], + ); + + let desc = admin + .registry_query_paginated( + ObjectType::Metric, + "timestamp", + false, + None, + Some(limit), + Some(desc_order[anchor_idx]), + Some(1), + false, + ) + .await + .object_ids() + .collect::>(); + let desc_size = std::cmp::min(limit, total - anchor_idx - 1); + assert_eq!( + desc, + desc_order[anchor_idx + 1..anchor_idx + 1 + desc_size], + "descending anchor={} offset=1 limit={limit}", + desc_order[anchor_idx], + ); + } + // Purge metrics and make sure they are gone test.server .metrics_store() diff --git a/tests/src/utils/registry.rs b/tests/src/utils/registry.rs index 1db627a5..21808d3b 100644 --- a/tests/src/utils/registry.rs +++ b/tests/src/utils/registry.rs @@ -144,6 +144,49 @@ impl Account { .await } + #[allow(clippy::too_many_arguments)] + pub async fn registry_query_paginated( + &self, + object: ObjectType, + sort_property: &str, + sort_ascending: bool, + position: Option, + limit: Option, + anchor: Option, + anchor_offset: Option, + calculate_total: bool, + ) -> JmapResponse { + let name = object.as_str(); + let mut args = serde_json::Map::new(); + args.insert("filter".into(), json!({})); + args.insert( + "sort".into(), + json!([{ "property": sort_property, "isAscending": sort_ascending }]), + ); + if let Some(p) = position { + args.insert("position".into(), json!(p)); + } + if let Some(l) = limit { + args.insert("limit".into(), json!(l)); + } + if let Some(a) = anchor { + args.insert("anchor".into(), json!(a.to_string())); + } + if let Some(ao) = anchor_offset { + args.insert("anchorOffset".into(), json!(ao)); + } + if calculate_total { + args.insert("calculateTotal".into(), json!(true)); + } + + self.jmap_method_calls(json!([[ + format!("x:{name}/query"), + Value::Object(args), + "0" + ]])) + .await + } + pub async fn registry_destroy( &self, object: ObjectType,