From 4fbae9218757247f919cb125e8a2a19a64d9f5e6 Mon Sep 17 00:00:00 2001 From: Abdul Andha Date: Fri, 24 Apr 2026 15:33:26 -0400 Subject: [PATCH] send after key on last page --- src/aggregation/bucket/composite/mod.rs | 85 ++++++++++++++-------- src/aggregation/intermediate_agg_result.rs | 32 ++++---- 2 files changed, 69 insertions(+), 48 deletions(-) diff --git a/src/aggregation/bucket/composite/mod.rs b/src/aggregation/bucket/composite/mod.rs index fe0a2a26e..0a6c30d3c 100644 --- a/src/aggregation/bucket/composite/mod.rs +++ b/src/aggregation/bucket/composite/mod.rs @@ -559,34 +559,30 @@ mod tests { page_size, agg_req, ); - if page_idx + 1 < page_count { - assert!( - res["my_composite"].get("after_key").is_some(), - "expected after_key on all but last page" - ); - after_key = Some(res["my_composite"]["after_key"].clone()); - } else if res["my_composite"].get("after_key").is_some() { - // currently we sometime have an after_key on the last page, - // check that the next "page" is empty - let agg_req_json = json!({ - "my_composite": { - "composite": { - "sources": composite_agg_sources, - "size": page_size, - "after": res["my_composite"]["after_key"].clone(), - } - } - }); - let agg_req: Aggregations = serde_json::from_value(agg_req_json).unwrap(); - let res = exec_request(agg_req.clone(), index).unwrap(); - assert_eq!( - res["my_composite"]["buckets"], - json!([]), - "expected no buckets when using after_key from last page, query: {:?}", - agg_req - ); - } + assert!( + res["my_composite"].get("after_key").is_some(), + "expected after_key on every non-empty page" + ); + after_key = Some(res["my_composite"]["after_key"].clone()); } + // Using the after_key from the last page must yield an empty page. + let agg_req_json = json!({ + "my_composite": { + "composite": { + "sources": composite_agg_sources, + "size": page_size, + "after": after_key, + } + } + }); + let agg_req: Aggregations = serde_json::from_value(agg_req_json).unwrap(); + let res = exec_request(agg_req.clone(), index).unwrap(); + assert_eq!( + res["my_composite"]["buckets"], + json!([]), + "expected no buckets when using after_key from last page, query: {:?}", + agg_req + ); } } @@ -711,8 +707,27 @@ mod tests { {"key": {"myterm": "terme"}, "doc_count": 1} ]) ); - assert!(res["my_composite"].get("after_key").is_none()); + // paginating past last page should be empty + let agg_req_json = json!({ + "my_composite": { + "composite": { + "sources": [ + {"myterm": {"terms": {"field": "string_id"}}} + ], + "size": 3, + "after": &res["my_composite"]["after_key"] + } + } + }); + let agg_req: Aggregations = serde_json::from_value(agg_req_json).unwrap(); + let res = exec_request(agg_req.clone(), &index).unwrap(); + assert_eq!( + res["my_composite"]["buckets"], + json!([]), + "expected no buckets when using after_key from last page, query: {:?}", + agg_req + ); Ok(()) } @@ -820,7 +835,10 @@ mod tests { {"key": {"myterm": "apple"}, "doc_count": 1} ]) ); - assert!(res["fruity_aggreg"].get("after_key").is_none()); + assert_eq!( + res["fruity_aggreg"]["after_key"], + json!({"myterm": "str:apple"}) + ); Ok(()) } @@ -1792,7 +1810,14 @@ mod tests { {"key": {"month": ms_timestamp_from_iso_str("2021-02-01T00:00:00Z"), "category": "books"}, "doc_count": 1}, ]), ); - assert!(res["my_composite"].get("after_key").is_none()); + let feb_2021_ns = ms_timestamp_from_iso_str("2021-02-01T00:00:00Z") * 1_000_000; + assert_eq!( + res["my_composite"]["after_key"], + json!({ + "month": format!("dt:{}", feb_2021_ns), + "category": "str:books" + }) + ); Ok(()) } diff --git a/src/aggregation/intermediate_agg_result.rs b/src/aggregation/intermediate_agg_result.rs index 6c6eeeb7e..3b6fcc076 100644 --- a/src/aggregation/intermediate_agg_result.rs +++ b/src/aggregation/intermediate_agg_result.rs @@ -1004,24 +1004,20 @@ impl IntermediateCompositeBucketResult { ) -> crate::Result { let trimmed_entry_vec = trim_composite_buckets(self.entries, &self.orders, self.target_size)?; - let after_key = if trimmed_entry_vec.len() == req.size as usize { - trimmed_entry_vec - .last() - .map(|bucket| { - let (intermediate_key, _entry) = bucket; - intermediate_key - .iter() - .enumerate() - .map(|(idx, intermediate_key)| { - let source = &req.sources[idx]; - (source.name().to_string(), intermediate_key.clone().into()) - }) - .collect() - }) - .unwrap() - } else { - FxHashMap::default() - }; + let after_key = trimmed_entry_vec + .last() + .map(|bucket| { + let (intermediate_key, _entry) = bucket; + intermediate_key + .iter() + .enumerate() + .map(|(idx, intermediate_key)| { + let source = &req.sources[idx]; + (source.name().to_string(), intermediate_key.clone().into()) + }) + .collect() + }) + .unwrap_or_default(); let buckets = trimmed_entry_vec .into_iter()