From 6c9c337853c2522fdf9c4d7692283120d4a942cc Mon Sep 17 00:00:00 2001 From: Darshit Chanpura Date: Thu, 13 Aug 2026 16:20:13 +0000 Subject: [PATCH 1/6] Preserve constant_score through filtered aliases When a search targets a filtered alias, DefaultSearchContext.buildFilteredQuery adds the alias filter by wrapping the user's query in a scoring BooleanQuery (bool[query MUST, aliasFilter FILTER]). If the user's query was a ConstantScoreQuery, nesting it inside that scoring boolean re-introduces scoring of the inner query and defeats Lucene's COMPLETE_NO_SCORES fast path, so the inner query's scorer is set up over the whole postings list. Re-wrap the filtered result in a ConstantScoreQuery when the main query was already constant-scored, keeping the outermost query non-scoring. This preserves both the no-scoring contract and the fast path end-to-end. Matters for access-control filters (e.g. document-level security composed with a filtered alias), where the combined query would otherwise be scored despite the caller having requested constant scoring. Verified by ConstantScoreFilteredAliasIT (hits score 1.0 through the alias, matching the backing index); IndexAliasesIT (29/29) confirms no regression for normal scored alias queries. Signed-off-by: Darshit Chanpura --- .../search/ConstantScoreFilteredAliasIT.java | 81 +++++++++++++++++++ .../search/DefaultSearchContext.java | 14 +++- 2 files changed, 94 insertions(+), 1 deletion(-) create mode 100644 server/src/internalClusterTest/java/org/opensearch/search/ConstantScoreFilteredAliasIT.java diff --git a/server/src/internalClusterTest/java/org/opensearch/search/ConstantScoreFilteredAliasIT.java b/server/src/internalClusterTest/java/org/opensearch/search/ConstantScoreFilteredAliasIT.java new file mode 100644 index 0000000000000..9886a554b18f8 --- /dev/null +++ b/server/src/internalClusterTest/java/org/opensearch/search/ConstantScoreFilteredAliasIT.java @@ -0,0 +1,81 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.opensearch.search; + +import org.opensearch.action.admin.indices.alias.IndicesAliasesRequest.AliasActions; +import org.opensearch.action.search.SearchResponse; +import org.opensearch.action.support.WriteRequest.RefreshPolicy; +import org.opensearch.common.settings.Settings; +import org.opensearch.index.query.QueryBuilders; +import org.opensearch.test.OpenSearchIntegTestCase; + +import static org.opensearch.test.hamcrest.OpenSearchAssertions.assertAcked; + +/** + * Verifies that a {@code constant_score} query keeps its non-scoring behavior when it is run through a + * filtered alias. Adding the alias filter wraps the user's query in a {@link org.apache.lucene.search.BooleanQuery} + * ({@code bool[query MUST, aliasFilter FILTER]}); without special handling this nests the + * {@link org.apache.lucene.search.ConstantScoreQuery} inside a scoring boolean, which both scores the + * hits (breaking the "no scoring" contract) and defeats Lucene's no-scoring fast path. + *

+ * {@link DefaultSearchContext#buildFilteredQuery} re-wraps the filtered result in a {@code ConstantScoreQuery} + * when the main query was already constant-scored, so the outermost query stays non-scoring. This asserts the + * observable consequence: every hit scores exactly 1.0 through the alias, identical to the same query against + * the backing index directly. + */ +public class ConstantScoreFilteredAliasIT extends OpenSearchIntegTestCase { + + private static final String INDEX = "docs"; + private static final String ALIAS = "docs_view"; + + private void buildCorpus() throws Exception { + assertAcked( + prepareCreate(INDEX).setMapping("dept", "type=keyword", "content", "type=text") + .setSettings(Settings.builder().put("index.number_of_shards", 1).put("index.number_of_replicas", 0)) + ); + // Enough docs carrying the term that BM25 scoring would produce clearly non-1.0, varying scores if the + // query were scored (different content lengths -> different norms). + for (int i = 0; i < 50; i++) { + client().prepareIndex(INDEX).setSource("dept", "cardiology", "content", ("alpha ".repeat(i % 5 + 1)) + "beta").get(); + } + for (int i = 0; i < 50; i++) { + client().prepareIndex(INDEX).setSource("dept", "oncology", "content", "alpha gamma " + i).get(); + } + client().prepareIndex(INDEX) + .setSource("dept", "cardiology", "content", "alpha probe") + .setRefreshPolicy(RefreshPolicy.IMMEDIATE) + .get(); + refresh(INDEX); + + AliasActions add = AliasActions.add().index(INDEX).alias(ALIAS).filter(QueryBuilders.termQuery("dept", "cardiology")); + assertAcked(client().admin().indices().prepareAliases().addAliasAction(add)); + } + + /** A constant_score query over "alpha" must score every hit exactly 1.0, on both the backing index and the + * filtered alias. Before the fix, the alias path scored hits (values != 1.0) because the alias filter + * nested the ConstantScoreQuery inside a scoring BooleanQuery. */ + public void testConstantScoreStaysConstantThroughFilteredAlias() throws Exception { + buildCorpus(); + + var constantScoreAlpha = QueryBuilders.constantScoreQuery(QueryBuilders.matchQuery("content", "alpha")); + + SearchResponse viaBacking = client().prepareSearch(INDEX).setQuery(constantScoreAlpha).setSize(200).get(); + SearchResponse viaAlias = client().prepareSearch(ALIAS).setQuery(constantScoreAlpha).setSize(200).get(); + + assertTrue("backing index returns hits", viaBacking.getHits().getTotalHits().value() > 0); + assertTrue("alias returns hits", viaAlias.getHits().getTotalHits().value() > 0); + + for (SearchHit hit : viaBacking.getHits().getHits()) { + assertEquals("constant_score on backing index must score 1.0", 1.0f, hit.getScore(), 0.0f); + } + for (SearchHit hit : viaAlias.getHits().getHits()) { + assertEquals("constant_score through filtered alias must also score 1.0", 1.0f, hit.getScore(), 0.0f); + } + } +} diff --git a/server/src/main/java/org/opensearch/search/DefaultSearchContext.java b/server/src/main/java/org/opensearch/search/DefaultSearchContext.java index d95321e8b0761..e50aba5fc434a 100644 --- a/server/src/main/java/org/opensearch/search/DefaultSearchContext.java +++ b/server/src/main/java/org/opensearch/search/DefaultSearchContext.java @@ -39,6 +39,7 @@ import org.apache.lucene.search.BoostQuery; import org.apache.lucene.search.Collector; import org.apache.lucene.search.CollectorManager; +import org.apache.lucene.search.ConstantScoreQuery; import org.apache.lucene.search.FieldDoc; import org.apache.lucene.search.MatchNoDocsQuery; import org.apache.lucene.search.Query; @@ -492,7 +493,18 @@ && new NestedHelper(mapperService()).mightMatchNestedDocs(query) for (Query filter : filters) { builder.add(filter, Occur.FILTER); } - return builder.build(); + BooleanQuery filtered = builder.build(); + // If the main query is already a ConstantScoreQuery, the caller has explicitly opted out of + // scoring. Adding the filters above nests that ConstantScoreQuery inside a scoring BooleanQuery, + // which defeats Lucene's COMPLETE_NO_SCORES fast path (the MUST clause's scorer is still set up + // over the whole postings list). Wrapping the filtered result back in a ConstantScoreQuery keeps + // the outermost query non-scoring, so the fast path is preserved end-to-end. This matters for + // access-control filters (e.g. document-level security composed with a filtered alias), where the + // combined query would otherwise be scored despite the caller having requested constant scoring. + if (query instanceof ConstantScoreQuery) { + return new ConstantScoreQuery(filtered); + } + return filtered; } } From b280eeaf017ca7d0ceab48c6e60d4922d79fc5ad Mon Sep 17 00:00:00 2001 From: Liyun Xiu Date: Thu, 20 Aug 2026 06:52:29 +0800 Subject: [PATCH 2/6] Apply http.max_header_size to HTTP/2 request headers (#22765) Signed-off-by: Liyun Xiu Co-authored-by: Andriy Redko --- .../netty4/Netty4HttpServerTransport.java | 17 +++++++-- .../Netty4HttpServerTransportTests.java | 38 +++++++++++++++++++ 2 files changed, 52 insertions(+), 3 deletions(-) diff --git a/modules/transport-netty4/src/main/java/org/opensearch/http/netty4/Netty4HttpServerTransport.java b/modules/transport-netty4/src/main/java/org/opensearch/http/netty4/Netty4HttpServerTransport.java index d125d8c17387d..dc9add81d680f 100644 --- a/modules/transport-netty4/src/main/java/org/opensearch/http/netty4/Netty4HttpServerTransport.java +++ b/modules/transport-netty4/src/main/java/org/opensearch/http/netty4/Netty4HttpServerTransport.java @@ -98,9 +98,11 @@ import io.netty.handler.codec.http.HttpServerUpgradeHandler.UpgradeCodecFactory; import io.netty.handler.codec.http2.CleartextHttp2ServerUpgradeHandler; import io.netty.handler.codec.http2.Http2CodecUtil; +import io.netty.handler.codec.http2.Http2FrameCodec; import io.netty.handler.codec.http2.Http2FrameCodecBuilder; import io.netty.handler.codec.http2.Http2MultiplexHandler; import io.netty.handler.codec.http2.Http2ServerUpgradeCodec; +import io.netty.handler.codec.http2.Http2Settings; import io.netty.handler.codec.http2.Http2StreamFrameToHttpObjectCodec; import io.netty.handler.logging.LogLevel; import io.netty.handler.logging.LoggingHandler; @@ -429,7 +431,7 @@ protected void configurePipeline(Channel ch) { public UpgradeCodec newUpgradeCodec(CharSequence protocol) { if (AsciiString.contentEquals(Http2CodecUtil.HTTP_UPGRADE_PROTOCOL_NAME, protocol)) { return new Http2ServerUpgradeCodec( - Http2FrameCodecBuilder.forServer().build(), + createHttp2FrameCodec(), new Http2MultiplexHandler(createHttp2ChannelInitializer(ch.pipeline())) ); } else { @@ -518,8 +520,17 @@ protected void configureDefaultHttpPipeline(ChannelPipeline pipeline) { } protected void configureDefaultHttp2Pipeline(ChannelPipeline pipeline) { - pipeline.addLast(Http2FrameCodecBuilder.forServer().build()) - .addLast(new Http2MultiplexHandler(createHttp2ChannelInitializer(pipeline))); + pipeline.addLast(createHttp2FrameCodec()).addLast(new Http2MultiplexHandler(createHttp2ChannelInitializer(pipeline))); + } + + /** + * Advertises {@code http.max_header_size} as SETTINGS_MAX_HEADER_LIST_SIZE, otherwise Netty's 8kb default + * caps HTTP/2 request headers regardless of the configured value. + */ + private Http2FrameCodec createHttp2FrameCodec() { + return Http2FrameCodecBuilder.forServer() + .initialSettings(Http2Settings.defaultSettings().maxHeaderListSize(handlingSettings.getMaxHeaderSize())) + .build(); } private ChannelInitializer createHttp2ChannelInitializerPriorKnowledge() { diff --git a/modules/transport-netty4/src/test/java/org/opensearch/http/netty4/Netty4HttpServerTransportTests.java b/modules/transport-netty4/src/test/java/org/opensearch/http/netty4/Netty4HttpServerTransportTests.java index 97b40b1fc8fd2..b2225cc6e6da5 100644 --- a/modules/transport-netty4/src/test/java/org/opensearch/http/netty4/Netty4HttpServerTransportTests.java +++ b/modules/transport-netty4/src/test/java/org/opensearch/http/netty4/Netty4HttpServerTransportTests.java @@ -490,6 +490,44 @@ protected void initChannel(SocketChannel ch) { } } + public void testLargeHeaderHttp2() throws InterruptedException { + final Settings settings = createBuilderWithPort().put(HttpTransportSettings.SETTING_HTTP_MAX_HEADER_SIZE.getKey(), "32kb").build(); + final String url = "/thing"; + final HttpServerTransport.Dispatcher dispatcher = dispatcherBuilderWithDefaults().withDispatchRequest( + (request, channel, threadContext) -> channel.sendResponse(new BytesRestResponse(OK, "done")) + ).build(); + + try ( + Netty4HttpServerTransport transport = new Netty4HttpServerTransport( + settings, + networkService, + bigArrays, + threadPool, + xContentRegistry(), + dispatcher, + clusterSettings, + new SharedGroupFactory(settings), + NoopTracer.INSTANCE + ) + ) { + transport.start(); + final TransportAddress remoteAddress = randomFrom(transport.boundAddress().boundAddresses()); + + try (Netty4HttpClient client = Netty4HttpClient.http2()) { + final FullHttpRequest request = new DefaultFullHttpRequest(HttpVersion.HTTP_1_1, HttpMethod.GET, url); + // Exceeds Netty's 8kb default for SETTINGS_MAX_HEADER_LIST_SIZE, but fits within http.max_header_size + request.headers().add("x-large-header", randomAlphaOfLength(16 * 1024)); + + final FullHttpResponse response = client.send(remoteAddress.address(), request); + try { + assertThat(response.status(), equalTo(HttpResponseStatus.OK)); + } finally { + response.release(); + } + } + } + } + private Settings createSettings() { return createBuilderWithPort().build(); } From e166ef6e72d5075b26b63c27f5c1cebf11c8d570 Mon Sep 17 00:00:00 2001 From: Peter Zhu Date: Wed, 19 Aug 2026 21:23:58 -0400 Subject: [PATCH 3/6] Bump 1password/load-secrets-action to v5.0.1 (OpenSearch) (#22778) Signed-off-by: Peter Zhu --- .github/workflows/publish-maven-snapshots.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/publish-maven-snapshots.yml b/.github/workflows/publish-maven-snapshots.yml index 045ce17b5122d..fd6e1386b2fa3 100644 --- a/.github/workflows/publish-maven-snapshots.yml +++ b/.github/workflows/publish-maven-snapshots.yml @@ -35,7 +35,7 @@ jobs: unzip protoc.zip -d $HOME/.local && rm -v protoc.zip && protoc --version - name: Load secret - uses: 1password/load-secrets-action@dafbe7cb03502b260e2b2893c753c352eee545bf # v3 + uses: 1password/load-secrets-action@70062d7a876d3eb6334754fa26efd2fbd90c32f2 # v5.0.1 with: # Export loaded secrets as environment variables export-env: true From ba2eb70ab5dd78f95c120efa3328bb7a689e5396 Mon Sep 17 00:00:00 2001 From: Somesh Gupta <35426854+aasom143@users.noreply.github.com> Date: Thu, 20 Aug 2026 13:57:29 +0530 Subject: [PATCH 4/6] Eliminate FFM round-trips via nextDoc piggyback on collectDocs (#22493) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit collectDocs now returns packed i64: upper 32 bits = next matching docId beyond maxDoc, lower 32 bits = wordsWritten. No new FFM calls needed. Rust consumers use nextDoc for three levels of optimization: 1. Full RG skip: nextDoc >= rgMax → no collectDocs call 2. Tightened range: nextDoc > rgMin → start from nextDoc, smaller bitset 3. No benefit: nextDoc <= rgMin → call as before Implemented in both SingleCollectorEvaluator (AtomicI32 across RGs) and CollectorLeafBitmaps/BitmapTree (per-leaf HashMap keyed by Arc identity). Signed-off-by: Somesh Gupta --- .../analytics/spi/FilterDelegationHandle.java | 8 +- .../rust/src/indexed_executor.rs | 6 +- .../rust/src/indexed_table/bool_tree.rs | 6 +- .../src/indexed_table/eval/bitmap_tree.rs | 395 +++++++++++++++++- .../rust/src/indexed_table/eval/mod.rs | 9 +- .../indexed_table/eval/single_collector.rs | 284 ++++++++++++- .../rust/src/indexed_table/ffm_callbacks.rs | 39 +- .../rust/src/indexed_table/index.rs | 25 +- .../tests_e2e/dynamic_filter_pushdown.rs | 9 +- .../tests_e2e/fuzz/delegation.rs | 9 +- .../indexed_table/tests_e2e/fuzz/harness.rs | 15 +- .../rust/src/indexed_table/tests_e2e/mod.rs | 15 +- .../indexed_table/tests_e2e/multi_segment.rs | 30 +- .../indexed_table/tests_e2e/null_columns.rs | 8 +- .../indexed_table/tests_e2e/page_pruning.rs | 13 +- .../tests_e2e/qtf_fetch_phase.rs | 6 +- .../tests_e2e/row_id_emission.rs | 24 +- .../tests_e2e/sort_reverse_row_id.rs | 15 +- .../indexfilter/FilterTreeCallbacks.java | 6 +- .../DelegationTaskTrackingTests.java | 2 +- .../indexfilter/IndexFilterCallbackTests.java | 4 +- .../lucene/LuceneFilterDelegationHandle.java | 17 +- 22 files changed, 836 insertions(+), 109 deletions(-) diff --git a/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/FilterDelegationHandle.java b/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/FilterDelegationHandle.java index 9fcffabd4fb6d..d9ccc413da62f 100644 --- a/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/FilterDelegationHandle.java +++ b/sandbox/libs/analytics-framework/src/main/java/org/opensearch/analytics/spi/FilterDelegationHandle.java @@ -62,13 +62,17 @@ public interface FilterDelegationHandle extends Closeable { *

Bit layout: word {@code i} contains matches for docs * {@code [minDoc + i*64, minDoc + (i+1)*64)}, LSB-first within each word. * + *

Return value packs two fields: upper 32 bits = next matching docId beyond + * maxDoc (or {@code Integer.MAX_VALUE} if exhausted), lower 32 bits = words written. + * Callers can use the nextDoc to skip subsequent RGs where {@code nextDoc >= rgMax}. + * * @param collectorKey key returned by {@link #createCollector(int, long, int, int)} * @param minDoc inclusive lower bound * @param maxDoc exclusive upper bound * @param out destination buffer; implementation writes up to {@code out.byteSize() / 8} words - * @return number of words written, or {@code -1} on error + * @return packed long: upper 32 = nextDoc, lower 32 = wordsWritten; or {@code -1} on error */ - int collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment out); + long collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment out); /** * Release resources for a collector. diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_executor.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_executor.rs index 1c662988f379f..b0f0e14df52c1 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_executor.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_executor.rs @@ -1360,9 +1360,9 @@ async unsafe fn execute_indexed_with_context_inner( let eval: Arc = Arc::new(TreeBitsetSource { tree: resolved, evaluator: Arc::new(BitmapTreeEvaluator), - leaves: Arc::new(CollectorLeafBitmaps { - ffm_collector_calls: stream_metrics.ffm_collector_calls.clone(), - }), + leaves: Arc::new(CollectorLeafBitmaps::new( + stream_metrics.ffm_collector_calls.clone(), + )), page_pruner: pruner, cost_predicate, cost_collector, diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/bool_tree.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/bool_tree.rs index 5ff047c8de8d8..8a3bba82a1147 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/bool_tree.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/bool_tree.rs @@ -31,6 +31,7 @@ use std::sync::Arc; use datafusion::physical_expr::PhysicalExpr; use super::index::RowGroupDocsCollector; +use crate::indexed_table::index::CollectDocsResult; /// A node in the boolean query tree (unresolved). #[derive(Debug, Clone)] @@ -463,6 +464,7 @@ pub fn residual_bool_to_physical_expr( #[cfg(test)] mod tests { use super::*; + use crate::indexed_table::index::CollectDocsResult; use crate::indexed_table::index::RowGroupDocsCollector; use datafusion::arrow::datatypes::{DataType, Field, Schema}; use datafusion::common::ScalarValue; @@ -473,8 +475,8 @@ mod tests { #[derive(Debug)] struct StubCollector(u8); impl RowGroupDocsCollector for StubCollector { - fn collect_packed_u64_bitset(&self, _: i32, _: i32) -> Result, String> { - Ok(vec![self.0 as u64]) + fn collect_packed_u64_bitset(&self, _: i32, _: i32) -> Result { + Ok(vec![self.0 as u64].into()) } } diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/bitmap_tree.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/bitmap_tree.rs index eafa3688d3dd6..e45152cfc2c9a 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/bitmap_tree.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/bitmap_tree.rs @@ -984,15 +984,42 @@ pub struct CollectorLeafBitmaps { /// round-trip to Java per Collector leaf per RG. `None` for tests /// that don't care about metrics. pub ffm_collector_calls: Option, + /// Per-leaf iterator position carried across row groups: maps a collector + /// leaf (keyed by its `Arc` pointer identity) to the next matching docId + /// returned by that leaf's last `collectDocs` call. A later RG whose whole + /// range sits below this value has no matches for the leaf and is skipped + /// without an FFM call. + /// + /// Why interior mutability at all: `leaf_bitmap` takes `&self` (the + /// `LeafBitmapSource` trait signature) but must update this map on every + /// RG, so it can't take `&mut self`. + /// + /// Why `Mutex` and not `RefCell`: `LeafBitmapSource: Send + Sync` and the + /// impl is held as `Arc`, so the whole struct must be + /// `Sync`. `RefCell` is `!Sync` and won't compile under that bound; `Mutex` + /// is the simplest `Sync` keyed cell. (`SingleCollectorEvaluator` uses a + /// single `AtomicI32` because it tracks exactly one collector; a tree has + /// many leaves, hence a keyed map.) + /// + /// The lock is uncontended in practice: within a partition, row groups are + /// prefetched sequentially, so no two threads call `leaf_bitmap` on the same + /// instance at once. The `Mutex` satisfies the `Sync` bound, not real + /// concurrent access. + leaf_to_next_doc_map: std::sync::Mutex>, } impl CollectorLeafBitmaps { - /// Construct a `CollectorLeafBitmaps` with no metrics. - pub fn without_metrics() -> Self { + pub fn new(ffm_collector_calls: Option) -> Self { Self { - ffm_collector_calls: None, + ffm_collector_calls, + leaf_to_next_doc_map: std::sync::Mutex::new(HashMap::new()), } } + + /// Construct a `CollectorLeafBitmaps` with no metrics. + pub fn without_metrics() -> Self { + Self::new(None) + } } impl LeafBitmapSource for CollectorLeafBitmaps { @@ -1008,7 +1035,26 @@ impl LeafBitmapSource for CollectorLeafBitmaps { return Err("CollectorLeafBitmaps: non-Collector node passed to leaf_bitmap".into()) } }; - // Use the narrowed call ranges if available (set by AND evaluator + + // `Arc` pointer identity is a stable per-query key: the collector Arc + // lives inside the resolved tree for the whole query, so its address + // never aliases another leaf mid-query. + let leaf_key = Arc::as_ptr(collector) as *const () as usize; + + // nextDoc from this leaf's previous collectDocs (i32::MIN = "no info yet"). + // Ranges are half-open [min, max): a nextDoc at or past a range's exclusive + // upper bound means no match in that range, so we compare with `>=`. + let last_next_doc = { + let map = self.leaf_to_next_doc_map.lock().unwrap(); + map.get(&leaf_key).copied().unwrap_or(i32::MIN) + }; + + // Whole RG is past the next match → no FFM call needed. + if last_next_doc >= ctx.max_doc { + return Ok(RoaringBitmap::new()); + } + + // Use the narrowed call ranges if available (set by the AND evaluator // after earlier children shrink the candidate set). Each range // produces one FFM call; results are merged into one bitmap in // min_doc-relative coordinates. @@ -1019,15 +1065,32 @@ impl LeafBitmapSource for CollectorLeafBitmaps { .unwrap_or_else(|| vec![(ctx.min_doc, ctx.max_doc)]); let mut result_bitmap = RoaringBitmap::new(); + let mut next_doc_out = last_next_doc; for (call_min, call_max) in &call_ranges { - let bitset = collector.collect_packed_u64_bitset(*call_min, *call_max)?; + // Sub-ranges are ascending; carry the freshest next_doc forward so a + // later sub-range skips/tightens on the position the iterator has + // already advanced to. Skip a sub-range entirely below the next match; + // otherwise tighten its lower bound so collectDocs skips the empty prefix. + if next_doc_out >= *call_max { + continue; + } + let effective_min = next_doc_out.max(*call_min); + let result = collector.collect_packed_u64_bitset(effective_min, *call_max)?; if let Some(ref c) = self.ffm_collector_calls { c.add(1); } - let offset = (*call_min - ctx.min_doc) as u32; - let num_docs = (*call_max - *call_min) as u32; + // Advance only forward — the iterator position is monotonic; guard + // against a stale/sentinel next_doc dragging it backward. + next_doc_out = next_doc_out.max(result.next_doc); + // Bitset is relative to effective_min; place it at the matching + // RG-relative offset so bit k maps to absolute doc effective_min + k. + let offset = (effective_min - ctx.min_doc) as u32; + let num_docs = (*call_max - effective_min) as u32; let bytes: &[u8] = unsafe { - std::slice::from_raw_parts(bitset.as_ptr() as *const u8, bitset.len() * 8) + std::slice::from_raw_parts( + result.words.as_ptr() as *const u8, + result.words.len() * 8, + ) }; let mut chunk = RoaringBitmap::from_lsb0_bytes(offset, bytes); let upper = offset + num_docs; @@ -1036,6 +1099,11 @@ impl LeafBitmapSource for CollectorLeafBitmaps { } result_bitmap |= chunk; } + + // Persist this leaf's advanced position for subsequent row groups. + let mut map = self.leaf_to_next_doc_map.lock().unwrap(); + map.insert(leaf_key, next_doc_out); + Ok(result_bitmap) } } @@ -1048,6 +1116,7 @@ impl LeafBitmapSource for CollectorLeafBitmaps { mod tests { use super::*; use crate::indexed_table::bool_tree::ResolvedNode; + use crate::indexed_table::index::CollectDocsResult; use crate::indexed_table::index::RowGroupDocsCollector; use datafusion::arrow::array::Int32Array; use datafusion::arrow::datatypes::{DataType, Field, Schema}; @@ -1120,8 +1189,12 @@ mod tests { #[derive(Debug)] struct Dummy; impl RowGroupDocsCollector for Dummy { - fn collect_packed_u64_bitset(&self, _: i32, _: i32) -> Result, String> { - Ok(vec![]) + fn collect_packed_u64_bitset( + &self, + _: i32, + _: i32, + ) -> Result { + Ok(vec![].into()) } } let _ = idx; @@ -1366,7 +1439,11 @@ mod tests { #[derive(Debug)] struct Poison; impl RowGroupDocsCollector for Poison { - fn collect_packed_u64_bitset(&self, _: i32, _: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + _: i32, + _: i32, + ) -> Result { unreachable!("Phase 2 must not call collect") } } @@ -2012,6 +2089,302 @@ mod tests { assert_eq!(result.candidates, bm(&[2, 3])); } + // ── CollectorLeafBitmaps next_doc skip/tighten tests ───────────── + + /// Mock collector that returns configurable docs and next_doc, and records + /// the (min_doc, max_doc) arguments it was called with. + #[derive(Debug)] + struct NextDocMockCollector { + /// Absolute doc IDs to include in the returned bitset. + docs: Vec, + /// The `next_doc` value to return from `collect_packed_u64_bitset`. + next_doc: i32, + /// Records each (min_doc, max_doc) invocation. + calls: std::sync::Mutex>, + } + + impl NextDocMockCollector { + fn new(docs: Vec, next_doc: i32) -> Self { + Self { + docs, + next_doc, + calls: std::sync::Mutex::new(Vec::new()), + } + } + + fn call_args(&self) -> Vec<(i32, i32)> { + self.calls.lock().unwrap().clone() + } + } + + impl RowGroupDocsCollector for NextDocMockCollector { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { + self.calls.lock().unwrap().push((min_doc, max_doc)); + // Mirror the real FfmSegmentCollector empty-range shortcut: an empty + // range yields no words and reports the scorer as exhausted. A skip + // check that lets an empty (min == max) call through would poison the + // stored next_doc to i32::MAX and drop later matches. + if max_doc <= min_doc { + return Ok(CollectDocsResult { + words: Vec::new(), + next_doc: i32::MAX, + }); + } + let num_docs = (max_doc - min_doc) as usize; + let num_words = num_docs.div_ceil(64); + let mut words = vec![0u64; num_words]; + for &d in &self.docs { + if d >= min_doc && d < max_doc { + let bit = (d - min_doc) as usize; + words[bit / 64] |= 1u64 << (bit % 64); + } + } + Ok(CollectDocsResult { + words, + next_doc: self.next_doc, + }) + } + } + + /// Helper: build a ResolvedNode::Collector wrapping a given Arc collector. + fn collector_node_from_arc(collector: Arc) -> ResolvedNode { + ResolvedNode::Collector { + provider_key: 1, + collector, + } + } + + /// Test 1: Per-leaf skip. + /// RG0 returns next_doc=500. RG1 covers [100, 200). + /// Since 500 > 200 (call_max), the sub-range is skipped entirely, + /// resulting in an empty bitmap for RG1. + #[test] + fn next_doc_skip_when_next_doc_exceeds_call_max() { + let mock = Arc::new(NextDocMockCollector::new(vec![110, 120, 130], 500)); + let node = collector_node_from_arc(mock.clone()); + let source = CollectorLeafBitmaps::without_metrics(); + + // RG0: covers [0, 100). Collector returns next_doc=500. + let ctx0 = RgEvalContext { + rg_idx: 0, + rg_first_row: 0, + rg_num_rows: 100, + min_doc: 0, + max_doc: 100, + cost_predicate: 1, + cost_collector: 10, + collector_call_ranges: None, + collector_strategy: super::super::CollectorCallStrategy::FullRange, + }; + let bm0 = source.leaf_bitmap(&node, 0, &ctx0).unwrap(); + // The mock returns docs in [0,100) that match — none do, so empty. + assert!(bm0.is_empty()); + // Verify the collector was called for RG0. + assert_eq!(mock.call_args().len(), 1); + + // RG1: covers [100, 200). Since last_next_doc=500 > 200 (call_max), + // the entire range is skipped — collector should NOT be called again. + let ctx1 = RgEvalContext { + rg_idx: 1, + rg_first_row: 100, + rg_num_rows: 100, + min_doc: 100, + max_doc: 200, + cost_predicate: 1, + cost_collector: 10, + collector_call_ranges: None, + collector_strategy: super::super::CollectorCallStrategy::FullRange, + }; + let bm1 = source.leaf_bitmap(&node, 0, &ctx1).unwrap(); + assert!( + bm1.is_empty(), + "RG1 should be empty because next_doc > max_doc" + ); + // Collector was NOT called for RG1 — still only 1 total call. + assert_eq!( + mock.call_args().len(), + 1, + "collector should not be called when next_doc > call_max" + ); + } + + /// Test 2: Per-leaf tighten. + /// RG0 returns next_doc=150. RG1 covers [100, 200). + /// effective_min = max(150, 100) = 150. The collector should be called + /// with min_doc=150, not 100. + #[test] + fn next_doc_tighten_effective_min() { + // Docs at 160, 170 — both within the tightened range [150, 200). + let mock = Arc::new(NextDocMockCollector::new(vec![160, 170], 150)); + let node = collector_node_from_arc(mock.clone()); + let source = CollectorLeafBitmaps::without_metrics(); + + // RG0: covers [0, 100). Returns next_doc=150. + let ctx0 = RgEvalContext { + rg_idx: 0, + rg_first_row: 0, + rg_num_rows: 100, + min_doc: 0, + max_doc: 100, + cost_predicate: 1, + cost_collector: 10, + collector_call_ranges: None, + collector_strategy: super::super::CollectorCallStrategy::FullRange, + }; + let _ = source.leaf_bitmap(&node, 0, &ctx0).unwrap(); + assert_eq!(mock.call_args(), vec![(0, 100)]); + + // RG1: covers [100, 200). last_next_doc=150, so effective_min=max(150,100)=150. + let ctx1 = RgEvalContext { + rg_idx: 1, + rg_first_row: 100, + rg_num_rows: 100, + min_doc: 100, + max_doc: 200, + cost_predicate: 1, + cost_collector: 10, + collector_call_ranges: None, + collector_strategy: super::super::CollectorCallStrategy::FullRange, + }; + let bm1 = source.leaf_bitmap(&node, 0, &ctx1).unwrap(); + + // Verify the collector was called with tightened min_doc=150, not 100. + let calls = mock.call_args(); + assert_eq!(calls.len(), 2); + assert_eq!(calls[1], (150, 200)); + + // Bitmap contains docs 160 and 170 at correct RG-relative positions. + // offset = (150 - 100) = 50. doc 160 at bit (160-150)=10 placed at 50+10=60. + assert!(bm1.contains(60), "doc 160 should be at position 60"); + assert!(bm1.contains(70), "doc 170 should be at position 70"); + assert_eq!(bm1.len(), 2); + } + + /// Test 3: Multiple leaves are independent. + /// Leaf A returns next_doc=500, leaf B returns next_doc=50. + /// For RG1 [100, 200): leaf A skips (500 > 200), leaf B does not (50 < 200). + #[test] + fn next_doc_multiple_leaves_independent() { + let mock_a = Arc::new(NextDocMockCollector::new(vec![], 500)); + let mock_b = Arc::new(NextDocMockCollector::new(vec![110, 120], 50)); + let node_a = collector_node_from_arc(mock_a.clone()); + let node_b = collector_node_from_arc(mock_b.clone()); + let source = CollectorLeafBitmaps::without_metrics(); + + // RG0: covers [0, 100). Both leaves are called. + let ctx0 = RgEvalContext { + rg_idx: 0, + rg_first_row: 0, + rg_num_rows: 100, + min_doc: 0, + max_doc: 100, + cost_predicate: 1, + cost_collector: 10, + collector_call_ranges: None, + collector_strategy: super::super::CollectorCallStrategy::FullRange, + }; + let _ = source.leaf_bitmap(&node_a, 0, &ctx0).unwrap(); + let _ = source.leaf_bitmap(&node_b, 1, &ctx0).unwrap(); + assert_eq!(mock_a.call_args().len(), 1); + assert_eq!(mock_b.call_args().len(), 1); + + // RG1: covers [100, 200). + let ctx1 = RgEvalContext { + rg_idx: 1, + rg_first_row: 100, + rg_num_rows: 100, + min_doc: 100, + max_doc: 200, + cost_predicate: 1, + cost_collector: 10, + collector_call_ranges: None, + collector_strategy: super::super::CollectorCallStrategy::FullRange, + }; + let bm_a = source.leaf_bitmap(&node_a, 0, &ctx1).unwrap(); + let bm_b = source.leaf_bitmap(&node_b, 1, &ctx1).unwrap(); + + // Leaf A: next_doc=500 > 200 (call_max) → skipped, empty bitmap. + assert!( + bm_a.is_empty(), + "leaf A should skip: next_doc=500 > max_doc=200" + ); + assert_eq!( + mock_a.call_args().len(), + 1, + "leaf A collector should not be called for RG1" + ); + + // Leaf B: next_doc=50 < 200 → not skipped, collector called. + assert!( + !bm_b.is_empty(), + "leaf B should not skip: next_doc=50 < max_doc=200" + ); + assert_eq!( + mock_b.call_args().len(), + 2, + "leaf B collector should be called for RG1" + ); + // Docs 110, 120 relative to min_doc=100 → bits 10, 20. + assert!(bm_b.contains(10)); + assert!(bm_b.contains(20)); + } + + /// Regression for the exclusive-boundary bug (Bharath's scenario): + /// 3 RGs of 100 docs each. A match sits at doc 200 — exactly the start of + /// RG2, i.e. exactly RG1's exclusive max_doc. + /// RG0 [0,100) match: doc 5, next_doc=200 + /// RG1 [100,200) no matches — must be skipped WITHOUT an empty collect call + /// RG2 [200,300) match: doc 200 — boundary doc, must be collected + /// With a `>` check, RG1 would tighten to an empty collect(200,200) → the + /// mock returns next_doc=i32::MAX → RG2 skipped → doc 200 silently dropped. + /// With `>=`, RG1 is skipped outright and doc 200 survives. + #[test] + fn next_doc_boundary_doc_not_dropped() { + let mock = Arc::new(NextDocMockCollector::new(vec![5, 200], 200)); + let node = collector_node_from_arc(mock.clone()); + let source = CollectorLeafBitmaps::without_metrics(); + + let mk_ctx = |rg_idx: usize, first: i64, min: i32, max: i32| RgEvalContext { + rg_idx, + rg_first_row: first, + rg_num_rows: 100, + min_doc: min, + max_doc: max, + cost_predicate: 1, + cost_collector: 10, + collector_call_ranges: None, + collector_strategy: super::super::CollectorCallStrategy::FullRange, + }; + + // RG0 [0,100): doc 5 matches, next_doc=200. + let bm0 = source.leaf_bitmap(&node, 0, &mk_ctx(0, 0, 0, 100)).unwrap(); + assert!(bm0.contains(5), "doc 5 should be collected in RG0"); + + // RG1 [100,200): next_doc=200 >= max_doc=200 → skipped, no collect call. + let bm1 = source + .leaf_bitmap(&node, 0, &mk_ctx(1, 100, 100, 200)) + .unwrap(); + assert!(bm1.is_empty(), "RG1 has no matches"); + assert_eq!( + mock.call_args().len(), + 1, + "RG1 must be skipped without a collect call (no empty-range poison)" + ); + + // RG2 [200,300): boundary doc 200 must be collected. + let bm2 = source + .leaf_bitmap(&node, 0, &mk_ctx(2, 200, 200, 300)) + .unwrap(); + assert!( + bm2.contains(0), + "boundary doc 200 (RG-relative pos 0) must not be dropped" + ); + } + /// When `rg_idx` IS in the reverse map and maps to a position where /// `rg_can_match[pos] == false`, the subtree is correctly pruned. #[test] diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/mod.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/mod.rs index edf77f49f32cc..7f998f3820929 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/mod.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/mod.rs @@ -724,6 +724,7 @@ fn precompute_collector_leaves<'a>( mod tests { use super::*; use crate::indexed_table::bool_tree::ResolvedNode; + use crate::indexed_table::index::CollectDocsResult; use crate::indexed_table::index::RowGroupDocsCollector; use crate::indexed_table::page_pruner::PagePruner; use datafusion::arrow::array::Int32Array; @@ -814,8 +815,12 @@ mod tests { #[derive(Debug)] struct Dummy; impl RowGroupDocsCollector for Dummy { - fn collect_packed_u64_bitset(&self, _: i32, _: i32) -> Result, String> { - Ok(vec![]) + fn collect_packed_u64_bitset( + &self, + _: i32, + _: i32, + ) -> Result { + Ok(vec![].into()) } } let source = TreeBitsetSource { diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/single_collector.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/single_collector.rs index 4eab3790cb4dc..4883d7abff180 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/single_collector.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/eval/single_collector.rs @@ -32,7 +32,7 @@ use roaring::RoaringBitmap; use super::{PrefetchedRg, RowGroupBitsetSource}; use crate::indexed_table::ffm_callbacks::{create_provider, FfmSegmentCollector, ProviderHandle}; -use crate::indexed_table::index::RowGroupDocsCollector; +use crate::indexed_table::index::{CollectDocsResult, RowGroupDocsCollector}; use crate::indexed_table::page_pruner::{PagePruneMetrics, PagePruner, StatsPruneTree}; use crate::indexed_table::row_selection::{ bitmap_to_packed_bits, packed_bits_to_boolean_array, row_selection_to_bitmap, PositionMap, @@ -186,6 +186,9 @@ pub struct SingleCollectorEvaluator { stats_prune_tree: Option>, /// Reverse map: absolute RG index → position in `rg_can_match` vectors. rg_index_to_pos: HashMap, + /// Next matching docId from the last collectDocs call. When next_doc >= rg.max_doc, + /// the RG can be skipped without an FFM call. Initialized to i32::MIN (no skip info). + last_next_doc: std::sync::atomic::AtomicI32, } /// Resources needed for per-RG bloom filter pruning. @@ -231,6 +234,7 @@ impl SingleCollectorEvaluator { bloom_config, stats_prune_tree, rg_index_to_pos, + last_next_doc: std::sync::atomic::AtomicI32::new(i32::MIN), } } } @@ -284,6 +288,21 @@ impl RowGroupBitsetSource for SingleCollectorEvaluator { } } + // Skip RG if the previous collectDocs told us the next match is beyond this RG. + let last_next = self + .last_next_doc + .load(std::sync::atomic::Ordering::Acquire); + // max_doc is exclusive, so nextDoc == max_doc also means "no match in this RG". + if last_next >= max_doc { + native_bridge_common::log_debug!( + "SingleCollector: skipping RG {} — nextDoc={} >= maxDoc={}", + rg.index, + last_next, + max_doc + ); + return Ok(None); + } + // Page-prune to discover which row ranges survive. let page_ranges: Option> = self.pruning_predicate.as_ref().and_then(|pp| { self.page_pruner @@ -357,10 +376,18 @@ impl RowGroupBitsetSource for SingleCollectorEvaluator { }; // Call collector for each range, merge into one RG-relative bitmap. + // Sub-ranges are ascending; carry the freshest next_doc forward so + // a later sub-range skips/tightens on the position the iterator has + // already advanced to, not the stale pre-loop value. let mut bm = RoaringBitmap::new(); + let mut next_doc_out = last_next; for (r_min, r_max) in &call_ranges { - let bitset = collector - .collect_packed_u64_bitset(*r_min, *r_max) + if next_doc_out >= *r_max { + continue; + } + let effective_min = next_doc_out.max(*r_min); + let result = collector + .collect_packed_u64_bitset(effective_min, *r_max) .map_err(|e| { format!( "collector.collect_packed_u64_bitset(rg={}, [{}, {})): {}", @@ -370,10 +397,16 @@ impl RowGroupBitsetSource for SingleCollectorEvaluator { if let Some(ref c) = self.ffm_collector_calls { c.add(1); } - let offset = (*r_min as i64 - rg.first_row) as u32; - let num_docs = (*r_max - *r_min) as u32; + // Advance only forward — the iterator position is monotonic; guard + // against a stale/sentinel next_doc dragging it backward. + next_doc_out = next_doc_out.max(result.next_doc); + let offset = (effective_min as i64 - rg.first_row) as u32; + let num_docs = (*r_max - effective_min) as u32; let bytes: &[u8] = unsafe { - std::slice::from_raw_parts(bitset.as_ptr() as *const u8, bitset.len() * 8) + std::slice::from_raw_parts( + result.words.as_ptr() as *const u8, + result.words.len() * 8, + ) }; let mut chunk = RoaringBitmap::from_lsb0_bytes(offset, bytes); let upper = offset.saturating_add(num_docs); @@ -382,6 +415,8 @@ impl RowGroupBitsetSource for SingleCollectorEvaluator { } bm |= chunk; } + self.last_next_doc + .store(next_doc_out, std::sync::atomic::Ordering::Release); // For FullRange and TightenOuterBounds, AND with page bitmap // to remove rows in dead pages that the collector scanned. @@ -477,7 +512,7 @@ impl RowGroupBitsetSource for SingleCollectorEvaluator { e ) })?; - let bitset = collector + let result = collector .collect_packed_u64_bitset(min_doc, max_doc) .map_err(|e| { format!( @@ -491,7 +526,10 @@ impl RowGroupBitsetSource for SingleCollectorEvaluator { let offset = (min_doc as i64 - rg.first_row) as u32; let num_docs = (max_doc - min_doc) as u32; let bytes: &[u8] = unsafe { - std::slice::from_raw_parts(bitset.as_ptr() as *const u8, bitset.len() * 8) + std::slice::from_raw_parts( + result.words.as_ptr() as *const u8, + result.words.len() * 8, + ) }; let mut peer_bm = RoaringBitmap::from_lsb0_bytes(offset, bytes); let upper = offset.saturating_add(num_docs); @@ -674,7 +712,7 @@ mod tests { &self, min_doc: i32, max_doc: i32, - ) -> Result, String> { + ) -> Result { let span = (max_doc - min_doc) as usize; let mut bitset = vec![0u64; (span + 63) / 64]; for &doc in &self.docs { @@ -683,7 +721,7 @@ mod tests { bitset[idx / 64] |= 1u64 << (idx % 64); } } - Ok(bitset) + Ok(bitset.into()) } } @@ -988,6 +1026,232 @@ mod tests { assert_eq!(got, vec![1u32, 5]); } + /// Mock collector that returns specific docs AND a configurable next_doc. + #[derive(Debug)] + struct NextDocCollector { + docs: Vec, + next_doc: i32, + } + + impl RowGroupDocsCollector for NextDocCollector { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { + let span = (max_doc - min_doc) as usize; + let mut words = vec![0u64; (span + 63) / 64]; + for &doc in &self.docs { + if doc >= min_doc && doc < max_doc { + let idx = (doc - min_doc) as usize; + words[idx / 64] |= 1u64 << (idx % 64); + } + } + Ok(CollectDocsResult { + words, + next_doc: self.next_doc, + }) + } + } + + #[test] + fn next_doc_skips_subsequent_rg() { + // RG0 [0,8): collector returns next_doc=20 (beyond RG1's max_doc=16) + // RG1 [8,16): should be skipped entirely + // RG2 [16,24): next_doc=20 < 24, should NOT be skipped (doc 20 is in range) + let collector = Arc::new(NextDocCollector { + docs: vec![2, 20], + next_doc: 20, + }) as Arc; + let pruner = minimal_page_pruner(); + let eval = SingleCollectorEvaluator::new( + Some(collector), + pruner, + None, + None, + None, + None, + CollectorCallStrategy::FullRange, + Arc::new(HashMap::new()), + 0, + Arc::new(FfmDelegatedBackendCollectorFactory), + 0, + None, + None, + HashMap::new(), + ); + + let rg0 = RowGroupInfo { + index: 0, + first_row: 0, + num_rows: 8, + }; + let rg1 = RowGroupInfo { + index: 1, + first_row: 8, + num_rows: 8, + }; + let rg2 = RowGroupInfo { + index: 2, + first_row: 16, + num_rows: 8, + }; + + // RG0: has doc 2 + let pf = eval.prefetch_rg(&rg0, 0, 8).unwrap().expect("has match"); + assert_eq!(pf.candidates.iter().collect::>(), vec![2u32]); + + // RG1: skipped (next_doc=20 > max_doc=16) + assert!(eval.prefetch_rg(&rg1, 8, 16).unwrap().is_none()); + + // RG2: NOT skipped (next_doc=20 < max_doc=24) + let pf = eval.prefetch_rg(&rg2, 16, 24).unwrap(); + assert!(pf.is_some()); + } + + #[test] + fn next_doc_tightens_min_doc() { + // Collector has doc at position 5. next_doc=5 means the iterator + // is at doc 5 after RG0. For RG1 [4,8), effective_min should be + // max(5, 4) = 5, not 4. + let collector = Arc::new(NextDocCollector { + docs: vec![1, 5], + next_doc: 5, + }) as Arc; + let pruner = minimal_page_pruner(); + let eval = SingleCollectorEvaluator::new( + Some(collector), + pruner, + None, + None, + None, + None, + CollectorCallStrategy::FullRange, + Arc::new(HashMap::new()), + 0, + Arc::new(FfmDelegatedBackendCollectorFactory), + 0, + None, + None, + HashMap::new(), + ); + + let rg0 = RowGroupInfo { + index: 0, + first_row: 0, + num_rows: 4, + }; + let rg1 = RowGroupInfo { + index: 1, + first_row: 4, + num_rows: 4, + }; + + // RG0 [0,4): doc 1 matches + let pf = eval.prefetch_rg(&rg0, 0, 4).unwrap().expect("has match"); + assert_eq!(pf.candidates.iter().collect::>(), vec![1u32]); + + // RG1 [4,8): doc 5 matches, effective_min tightened to 5 + let pf = eval.prefetch_rg(&rg1, 4, 8).unwrap().expect("has match"); + assert_eq!(pf.candidates.iter().collect::>(), vec![1u32]); // doc 5 is at RG-relative pos 1 + } + + #[test] + fn next_doc_max_value_skips_all_remaining() { + // next_doc = i32::MAX means scorer exhausted — all subsequent RGs skipped + let collector = Arc::new(NextDocCollector { + docs: vec![0], + next_doc: i32::MAX, + }) as Arc; + let pruner = minimal_page_pruner(); + let eval = SingleCollectorEvaluator::new( + Some(collector), + pruner, + None, + None, + None, + None, + CollectorCallStrategy::FullRange, + Arc::new(HashMap::new()), + 0, + Arc::new(FfmDelegatedBackendCollectorFactory), + 0, + None, + None, + HashMap::new(), + ); + + let rg0 = RowGroupInfo { + index: 0, + first_row: 0, + num_rows: 8, + }; + let rg1 = RowGroupInfo { + index: 1, + first_row: 8, + num_rows: 8, + }; + + // RG0: has doc 0 + assert!(eval.prefetch_rg(&rg0, 0, 8).unwrap().is_some()); + + // RG1: skipped (next_doc=MAX > any max_doc) + assert!(eval.prefetch_rg(&rg1, 8, 16).unwrap().is_none()); + } + + #[test] + fn next_doc_at_rg_boundary_is_not_dropped() { + // Regression for the exclusive-boundary bug: a match sitting exactly at + // an RG's start (== previous RG's exclusive max_doc) must NOT be dropped. + // docs: 5 (in RG0) and 8 (start of RG1, == RG0's max_doc=8). + let collector = Arc::new(NextDocCollector { + docs: vec![5, 8], + next_doc: 8, + }) as Arc; + let pruner = minimal_page_pruner(); + let eval = SingleCollectorEvaluator::new( + Some(collector), + pruner, + None, + None, + None, + None, + CollectorCallStrategy::FullRange, + Arc::new(HashMap::new()), + 0, + Arc::new(FfmDelegatedBackendCollectorFactory), + 0, + None, + None, + HashMap::new(), + ); + + let rg0 = RowGroupInfo { + index: 0, + first_row: 0, + num_rows: 8, + }; + let rg1 = RowGroupInfo { + index: 1, + first_row: 8, + num_rows: 8, + }; + + // RG0 [0,8): doc 5 matches, next_doc=8 (boundary). + let pf0 = eval + .prefetch_rg(&rg0, 0, 8) + .unwrap() + .expect("RG0 has doc 5"); + assert_eq!(pf0.candidates.iter().collect::>(), vec![5u32]); + + // RG1 [8,16): next_doc=8 == min_doc, NOT skipped. Doc 8 must be collected. + let pf1 = eval + .prefetch_rg(&rg1, 8, 16) + .unwrap() + .expect("RG1 has boundary doc 8"); + assert_eq!(pf1.candidates.iter().collect::>(), vec![0u32]); // doc 8 at RG-relative pos 0 + } + // Keep the `fmt` import used #[allow(dead_code)] fn _use(_: &dyn fmt::Debug) {} diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/ffm_callbacks.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/ffm_callbacks.rs index d661c593f60be..8b161b85d1eae 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/ffm_callbacks.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/ffm_callbacks.rs @@ -26,7 +26,7 @@ use std::sync::atomic::{AtomicPtr, Ordering}; -use super::index::RowGroupDocsCollector; +use super::index::{CollectDocsResult, RowGroupDocsCollector}; // ── Callback signatures ─────────────────────────────────────────────── @@ -209,15 +209,22 @@ impl FfmSegmentCollector { } impl RowGroupDocsCollector for FfmSegmentCollector { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { if max_doc <= min_doc { - return Ok(Vec::new()); + return Ok(CollectDocsResult { + words: Vec::new(), + next_doc: i32::MAX, + }); } let span = (max_doc - min_doc) as usize; let word_count = span.div_ceil(64); let mut buf = vec![0u64; word_count]; let collect_fn = load_collect_docs()?; - let n = unsafe { + let packed = unsafe { collect_fn( self.context_id, self.key, @@ -227,27 +234,27 @@ impl RowGroupDocsCollector for FfmSegmentCollector { word_count as i64, ) }; - if n < 0 { + if packed < 0 { return Err(format!( "collectDocs(context_id={}, key={}) failed: {}", - self.context_id, self.key, n + self.context_id, self.key, packed )); } - // Defensive: the Java callback is contracted to return - // `wordsWritten <= outWordCap`. If it lied, the buffer already - // overflowed, but truncating won't recover the clobbered heap. - // Detect the violation and fail loudly so the Java callback bug - // is surfaced before downstream code consumes the tainted bitset. - let n = n as usize; - if n > word_count { + let packed_u = packed as u64; + let words_written = (packed_u & 0xFFFF_FFFF) as usize; + let next_doc = (packed_u >> 32) as i32; + if words_written > word_count { return Err(format!( "collectDocs(context_id={}, key={}) reported wordsWritten={} > capacity={}; \ callback contract violated (possible heap overflow)", - self.context_id, self.key, n, word_count, + self.context_id, self.key, words_written, word_count, )); } - buf.truncate(n); - Ok(buf) + buf.truncate(words_written); + Ok(CollectDocsResult { + words: buf, + next_doc, + }) } } diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/index.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/index.rs index 0e0283c3ed378..9eb686e9cfc92 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/index.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/index.rs @@ -45,8 +45,31 @@ use std::sync::Arc; /// (zero-length bitset). This is a no-op case and must not error. Callers /// rely on this — e.g. `IndexedStream` skips filter-bitset fetch on empty /// row groups by calling with `max_doc == min_doc`. +/// Result of a `collect_packed_u64_bitset` call. +#[derive(Debug, Clone)] +pub struct CollectDocsResult { + /// The packed u64 bitset of matching docs within the requested range. + pub words: Vec, + /// The next matching doc ID beyond `max_doc`, or `i32::MAX` if exhausted. + /// Callers can use this to skip subsequent RGs where `next_doc >= rg_max`. + pub next_doc: i32, +} + +impl From> for CollectDocsResult { + fn from(words: Vec) -> Self { + CollectDocsResult { + words, + next_doc: i32::MIN, + } + } +} + pub trait RowGroupDocsCollector: Send + Sync + Debug { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String>; + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result; } /// A searcher scoped to a single shard (index), created once per query. diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/dynamic_filter_pushdown.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/dynamic_filter_pushdown.rs index 61935d5409be5..61a24114bb4ac 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/dynamic_filter_pushdown.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/dynamic_filter_pushdown.rs @@ -32,6 +32,7 @@ use super::super::index::RowGroupDocsCollector; use super::super::page_pruner::PagePruner; use super::super::stream::{FilterStrategy, RowGroupInfo}; use super::super::table_provider::{IndexedTableConfig, IndexedTableProvider, SegmentFileInfo}; +use crate::indexed_table::index::CollectDocsResult; /// 16 rows, `price` = 15..0 **descending** in file order, 4 rows per row group → /// RG ranges (by row position) [15..12], [11..8], [7..4], [3..0]. The scan reads @@ -78,13 +79,17 @@ fn write_fixture() -> (NamedTempFile, SchemaRef) { struct MatchAllCollector; impl RowGroupDocsCollector for MatchAllCollector { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { let span = (max_doc - min_doc).max(0) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for rel in 0..span { out[rel / 64] |= 1u64 << (rel % 64); } - Ok(out) + Ok(out.into()) } } diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/fuzz/delegation.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/fuzz/delegation.rs index 6d06743867269..fe0aa9e606639 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/fuzz/delegation.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/fuzz/delegation.rs @@ -49,6 +49,7 @@ use crate::indexed_table::eval::single_collector::{ }; use crate::indexed_table::eval::RowGroupBitsetSource; use crate::indexed_table::ffm_callbacks::ProviderHandle; +use crate::indexed_table::index::CollectDocsResult; use crate::indexed_table::index::RowGroupDocsCollector; use crate::indexed_table::page_pruner::PagePruner; use crate::indexed_table::table_provider::{ @@ -212,7 +213,11 @@ struct StaticBitsetCollector { } impl RowGroupDocsCollector for StaticBitsetCollector { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for &doc in &self.matching { @@ -221,7 +226,7 @@ impl RowGroupDocsCollector for StaticBitsetCollector { out[rel / 64] |= 1u64 << (rel % 64); } } - Ok(out) + Ok(out.into()) } } diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/fuzz/harness.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/fuzz/harness.rs index bdd4598bf839b..4b501948ce8fa 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/fuzz/harness.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/fuzz/harness.rs @@ -29,6 +29,7 @@ use crate::indexed_table::bool_tree::BoolNode; use crate::indexed_table::eval::bitmap_tree::{BitmapTreeEvaluator, CollectorLeafBitmaps}; use crate::indexed_table::eval::single_collector::SingleCollectorEvaluator; use crate::indexed_table::eval::{RowGroupBitsetSource, TreeBitsetSource}; +use crate::indexed_table::index::CollectDocsResult; use crate::indexed_table::index::RowGroupDocsCollector; use crate::indexed_table::page_pruner::PagePruner; use crate::indexed_table::stream::{FilterStrategy, RowGroupInfo}; @@ -46,7 +47,11 @@ struct MockCollector { } impl RowGroupDocsCollector for MockCollector { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for &doc in &self.matching { @@ -55,7 +60,7 @@ impl RowGroupDocsCollector for MockCollector { out[rel / 64] |= 1u64 << (rel % 64); } } - Ok(out) + Ok(out.into()) } } @@ -259,9 +264,9 @@ pub(in crate::indexed_table::tests_e2e) async fn execute_tree_with_plan_pushdown let eval: Arc = Arc::new(TreeBitsetSource { tree: Arc::new(resolved), evaluator: Arc::new(BitmapTreeEvaluator), - leaves: Arc::new(CollectorLeafBitmaps { - ffm_collector_calls: stream_metrics.ffm_collector_calls.clone(), - }), + leaves: Arc::new(CollectorLeafBitmaps::new( + stream_metrics.ffm_collector_calls.clone(), + )), page_pruner: pruner, cost_predicate: 1, cost_collector: 10, diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/mod.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/mod.rs index 968695a8cd0ae..d185cf9fc27b7 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/mod.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/mod.rs @@ -37,6 +37,7 @@ use super::index::RowGroupDocsCollector; use super::page_pruner::PagePruner; use super::stream::{FilterStrategy, RowGroupInfo}; use super::table_provider::{IndexedTableConfig, IndexedTableProvider, SegmentFileInfo}; +use crate::indexed_table::index::CollectDocsResult; mod boolean_algebra; mod constant_predicate; @@ -152,7 +153,11 @@ struct MockCollector { } impl RowGroupDocsCollector for MockCollector { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for &doc in &self.matching { @@ -161,7 +166,7 @@ impl RowGroupDocsCollector for MockCollector { out[rel / 64] |= 1u64 << (rel % 64); } } - Ok(out) + Ok(out.into()) } } @@ -264,9 +269,9 @@ async fn run_tree_and_plan( tree: Arc::new(resolved), evaluator: Arc::new(BitmapTreeEvaluator), leaves: Arc::new( - crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps { - ffm_collector_calls: _stream_metrics.ffm_collector_calls.clone(), - }, + crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps::new( + _stream_metrics.ffm_collector_calls.clone(), + ), ), page_pruner: pruner, cost_predicate: 1, diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/multi_segment.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/multi_segment.rs index c29ab8f456441..d8d6eaabde52c 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/multi_segment.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/multi_segment.rs @@ -74,7 +74,11 @@ struct PerSegmentCollector { } impl RowGroupDocsCollector for PerSegmentCollector { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for &doc in &self.matching { @@ -83,7 +87,7 @@ impl RowGroupDocsCollector for PerSegmentCollector { out[rel / 64] |= 1u64 << (rel % 64); } } - Ok(out) + Ok(out.into()) } } @@ -272,7 +276,11 @@ struct ConcurrencyWitnessCollector { } impl RowGroupDocsCollector for ConcurrencyWitnessCollector { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { use std::sync::atomic::Ordering; let cur = self.in_flight.fetch_add(1, Ordering::SeqCst) + 1; // Update high-water-mark with a CAS loop. @@ -1016,13 +1024,13 @@ async fn run_wide_segments( &self, min_doc: i32, max_doc: i32, - ) -> Result, String> { + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for i in 0..span { out[i / 64] |= 1u64 << (i % 64); } - Ok(out) + Ok(out.into()) } } @@ -1365,13 +1373,13 @@ async fn run_wide_segments_with_stats_pruning( &self, min_doc: i32, max_doc: i32, - ) -> Result, String> { + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for i in 0..span { out[i / 64] |= 1u64 << (i % 64); } - Ok(out) + Ok(out.into()) } } @@ -1760,13 +1768,13 @@ async fn stats_prune_direct_prefetch_asserts_pruning_and_empty_bitsets() { &self, min_doc: i32, max_doc: i32, - ) -> Result, String> { + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for i in 0..span { out[i / 64] |= 1u64 << (i % 64); } - Ok(out) + Ok(out.into()) } } @@ -1961,13 +1969,13 @@ async fn stats_prune_asserts_empty_collector_bitset_in_pruned_subtree() { &self, min_doc: i32, max_doc: i32, - ) -> Result, String> { + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for i in 0..span { out[i / 64] |= 1u64 << (i % 64); } - Ok(out) + Ok(out.into()) } } diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/null_columns.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/null_columns.rs index 559cd0e87aba4..4ef0ddaf3dc4d 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/null_columns.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/null_columns.rs @@ -106,7 +106,11 @@ struct RgScopedCollector { } impl RowGroupDocsCollector for RgScopedCollector { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for &doc in &self.matching_rows { @@ -115,7 +119,7 @@ impl RowGroupDocsCollector for RgScopedCollector { out[rel / 64] |= 1u64 << (rel % 64); } } - Ok(out) + Ok(out.into()) } } diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/page_pruning.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/page_pruning.rs index a4bec23e41685..8a96ffca2c09d 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/page_pruning.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/page_pruning.rs @@ -63,6 +63,7 @@ use crate::indexed_table::eval::single_collector::SingleCollectorEvaluator; use crate::indexed_table::eval::{ CollectorCallStrategy, RgEvalContext, RowGroupBitsetSource, TreeBitsetSource, }; +use crate::indexed_table::index::CollectDocsResult; use crate::indexed_table::index::RowGroupDocsCollector; use crate::indexed_table::page_pruner::{build_pruning_predicate, PagePruner}; use crate::indexed_table::stream::{FilterStrategy, RowGroupInfo}; @@ -153,7 +154,11 @@ struct MockCollector { } impl RowGroupDocsCollector for MockCollector { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for &doc in &self.docs { @@ -162,7 +167,7 @@ impl RowGroupDocsCollector for MockCollector { out[rel / 64] |= 1u64 << (rel % 64); } } - Ok(out) + Ok(out.into()) } } @@ -320,9 +325,7 @@ async fn run_bitmap_tree(tree: BoolNode) -> (Vec, Arc) { let eval: Arc = Arc::new(TreeBitsetSource { tree: Arc::new(resolved), evaluator: Arc::new(BitmapTreeEvaluator), - leaves: Arc::new(CollectorLeafBitmaps { - ffm_collector_calls: sm.ffm_collector_calls.clone(), - }), + leaves: Arc::new(CollectorLeafBitmaps::new(sm.ffm_collector_calls.clone())), page_pruner: pruner, cost_predicate: 1, cost_collector: 10, diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/qtf_fetch_phase.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/qtf_fetch_phase.rs index b678cdf2fff53..e8d76a411336a 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/qtf_fetch_phase.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/qtf_fetch_phase.rs @@ -90,9 +90,9 @@ async fn query_phase(tree: BoolNode) -> Vec { tree: Arc::new(resolved), evaluator: Arc::new(BitmapTreeEvaluator), leaves: Arc::new( - crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps { - ffm_collector_calls: _stream_metrics.ffm_collector_calls.clone(), - }, + crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps::new( + _stream_metrics.ffm_collector_calls.clone(), + ), ), page_pruner: pruner, cost_predicate: 1, diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/row_id_emission.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/row_id_emission.rs index 5b77b49d7553e..74ca72ebb2e41 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/row_id_emission.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/row_id_emission.rs @@ -79,9 +79,9 @@ async fn run_tree_row_ids(tree: BoolNode) -> Vec { tree: Arc::new(resolved), evaluator: Arc::new(BitmapTreeEvaluator), leaves: Arc::new( - crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps { - ffm_collector_calls: _stream_metrics.ffm_collector_calls.clone(), - }, + crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps::new( + _stream_metrics.ffm_collector_calls.clone(), + ), ), page_pruner: pruner, cost_predicate: 1, @@ -256,9 +256,9 @@ async fn run_tree_row_ids_with_global_base(tree: BoolNode, global_base: u64) -> tree: Arc::new(resolved), evaluator: Arc::new(BitmapTreeEvaluator), leaves: Arc::new( - crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps { - ffm_collector_calls: _stream_metrics.ffm_collector_calls.clone(), - }, + crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps::new( + _stream_metrics.ffm_collector_calls.clone(), + ), ), page_pruner: pruner, cost_predicate: 1, @@ -543,9 +543,9 @@ async fn test_row_id_with_data_columns() { tree: Arc::new(resolved), evaluator: Arc::new(BitmapTreeEvaluator), leaves: Arc::new( - crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps { - ffm_collector_calls: _stream_metrics.ffm_collector_calls.clone(), - }, + crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps::new( + _stream_metrics.ffm_collector_calls.clone(), + ), ), page_pruner: pruner, cost_predicate: 1, @@ -815,9 +815,9 @@ async fn run_two_segments_row_ids(tree: BoolNode) -> Vec { tree: Arc::new(resolved), evaluator: Arc::new(BitmapTreeEvaluator), leaves: Arc::new( - crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps { - ffm_collector_calls: _stream_metrics.ffm_collector_calls.clone(), - }, + crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps::new( + _stream_metrics.ffm_collector_calls.clone(), + ), ), page_pruner: pruner, cost_predicate: 1, diff --git a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/sort_reverse_row_id.rs b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/sort_reverse_row_id.rs index 715d3bbb07567..d2d4eeaa78611 100644 --- a/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/sort_reverse_row_id.rs +++ b/sandbox/plugins/analytics-backend-datafusion/rust/src/indexed_table/tests_e2e/sort_reverse_row_id.rs @@ -38,6 +38,7 @@ use crate::indexed_executor::reverse_segment_iteration_order; use crate::indexed_table::bool_tree::BoolNode; use crate::indexed_table::eval::bitmap_tree::BitmapTreeEvaluator; use crate::indexed_table::eval::{RowGroupBitsetSource, TreeBitsetSource}; +use crate::indexed_table::index::CollectDocsResult; use crate::indexed_table::index::RowGroupDocsCollector; use crate::indexed_table::page_pruner::PagePruner; use crate::indexed_table::stream::{FilterStrategy, RowGroupInfo}; @@ -102,7 +103,11 @@ struct LocalOffsetCollector { } impl RowGroupDocsCollector for LocalOffsetCollector { - fn collect_packed_u64_bitset(&self, min_doc: i32, max_doc: i32) -> Result, String> { + fn collect_packed_u64_bitset( + &self, + min_doc: i32, + max_doc: i32, + ) -> Result { let span = (max_doc - min_doc) as usize; let mut out = vec![0u64; span.div_ceil(64)]; for &doc in &self.matching { @@ -111,7 +116,7 @@ impl RowGroupDocsCollector for LocalOffsetCollector { out[rel / 64] |= 1u64 << (rel % 64); } } - Ok(out) + Ok(out.into()) } } @@ -202,9 +207,9 @@ async fn collect_row_ids( tree: Arc::new(resolved), evaluator: Arc::new(BitmapTreeEvaluator), leaves: Arc::new( - crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps { - ffm_collector_calls: _stream_metrics.ffm_collector_calls.clone(), - }, + crate::indexed_table::eval::bitmap_tree::CollectorLeafBitmaps::new( + _stream_metrics.ffm_collector_calls.clone(), + ), ), page_pruner: pruner, cost_predicate: 1, diff --git a/sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/indexfilter/FilterTreeCallbacks.java b/sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/indexfilter/FilterTreeCallbacks.java index a3c86cc019333..1e3fe8f0bce43 100644 --- a/sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/indexfilter/FilterTreeCallbacks.java +++ b/sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/indexfilter/FilterTreeCallbacks.java @@ -223,7 +223,7 @@ public static int createCollector(long contextId, int providerKey, long writerGe } /** - * {@code collectDocs(contextId, collectorKey, minDoc, maxDoc, outPtr, outWordCap) -> wordsWritten|-1}. + * {@code collectDocs(contextId, collectorKey, minDoc, maxDoc, outPtr, outWordCap) -> packed(nextDoc|wordsWritten)|-1}. */ public static long collectDocs(long contextId, int collectorKey, int minDoc, int maxDoc, MemorySegment outPtr, long outWordCap) { long tid = trackStart(contextId); @@ -239,8 +239,8 @@ public static long collectDocs(long contextId, int collectorKey, int minDoc, int } int maxWords = (int) Math.min(outWordCap, (long) Integer.MAX_VALUE); MemorySegment view = outPtr.reinterpret((long) maxWords * Long.BYTES); - int wordsWritten = handle.collectDocs(collectorKey, minDoc, maxDoc, view); - return (wordsWritten < 0) ? -1L : wordsWritten; + long result = handle.collectDocs(collectorKey, minDoc, maxDoc, view); + return (result < 0) ? -1L : result; } catch (AssertionError e) { throw e; } catch (Throwable throwable) { diff --git a/sandbox/plugins/analytics-backend-datafusion/src/test/java/org/opensearch/be/datafusion/indexfilter/DelegationTaskTrackingTests.java b/sandbox/plugins/analytics-backend-datafusion/src/test/java/org/opensearch/be/datafusion/indexfilter/DelegationTaskTrackingTests.java index 0f048de49c37b..434a6ccb7bcb3 100644 --- a/sandbox/plugins/analytics-backend-datafusion/src/test/java/org/opensearch/be/datafusion/indexfilter/DelegationTaskTrackingTests.java +++ b/sandbox/plugins/analytics-backend-datafusion/src/test/java/org/opensearch/be/datafusion/indexfilter/DelegationTaskTrackingTests.java @@ -324,7 +324,7 @@ public int createCollector(int providerKey, long writerGeneration, int minDoc, i } @Override - public int collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment out) { + public long collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment out) { int n = Math.min(words.length, (int) (out.byteSize() / Long.BYTES)); for (int i = 0; i < n; i++) out.setAtIndex(ValueLayout.JAVA_LONG, i, words[i]); diff --git a/sandbox/plugins/analytics-backend-datafusion/src/test/java/org/opensearch/be/datafusion/indexfilter/IndexFilterCallbackTests.java b/sandbox/plugins/analytics-backend-datafusion/src/test/java/org/opensearch/be/datafusion/indexfilter/IndexFilterCallbackTests.java index 58dc6501f3abf..48efe41d530ad 100644 --- a/sandbox/plugins/analytics-backend-datafusion/src/test/java/org/opensearch/be/datafusion/indexfilter/IndexFilterCallbackTests.java +++ b/sandbox/plugins/analytics-backend-datafusion/src/test/java/org/opensearch/be/datafusion/indexfilter/IndexFilterCallbackTests.java @@ -108,7 +108,7 @@ public int createCollector(int providerKey, long writerGeneration, int minDoc, i } @Override - public int collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment out) { + public long collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment out) { return -1; } @@ -165,7 +165,7 @@ public int createCollector(int providerKey, long writerGeneration, int minDoc, i } @Override - public int collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment out) { + public long collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment out) { this.lastCollectorKey = collectorKey; int wordCount = Math.min(cannedWords.length, (int) (out.byteSize() / Long.BYTES)); for (int i = 0; i < wordCount; i++) { diff --git a/sandbox/plugins/analytics-backend-lucene/src/main/java/org/opensearch/be/lucene/LuceneFilterDelegationHandle.java b/sandbox/plugins/analytics-backend-lucene/src/main/java/org/opensearch/be/lucene/LuceneFilterDelegationHandle.java index 6a829cf52e5bf..c926b7f6ec529 100644 --- a/sandbox/plugins/analytics-backend-lucene/src/main/java/org/opensearch/be/lucene/LuceneFilterDelegationHandle.java +++ b/sandbox/plugins/analytics-backend-lucene/src/main/java/org/opensearch/be/lucene/LuceneFilterDelegationHandle.java @@ -223,7 +223,7 @@ public boolean isCancelled() { } @Override - public int collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment out) { + public long collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment out) { ScorerHandle handle = scorersByCollectorKey.get(collectorKey); if (handle == null) { return -1; @@ -233,6 +233,7 @@ public int collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment o } int span = maxDoc - minDoc; FixedBitSet bits = new FixedBitSet(span); + int nextDoc = Integer.MAX_VALUE; if (handle.scorer != null) { int scanFrom = Math.max(minDoc, handle.partitionMinDoc); @@ -252,9 +253,16 @@ public int collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment o } handle.currentDoc = docId; } + nextDoc = handle.currentDoc; } catch (IOException exception) { LOGGER.warn("IOException during collectDocs, returning partial bitset", exception); + // Iteration is only partial — don't signal exhaustion (MAX_VALUE), + // which would make callers skip all subsequent RGs for this leaf. + // Report maxDoc conservatively so later RGs are still probed. + nextDoc = maxDoc; } + } else { + nextDoc = handle.currentDoc; } } @@ -263,15 +271,16 @@ public int collectDocs(int collectorKey, int minDoc, int maxDoc, MemorySegment o MemorySegment.copy(words, 0, out, ValueLayout.JAVA_LONG, 0, wordCount); if (LOGGER.isDebugEnabled()) { LOGGER.debug( - "[scf] collectDocs collectorKey={} range=[{},{}) → cardinality={} words={}", + "[scf] collectDocs collectorKey={} range=[{},{}) → cardinality={} words={} nextDoc={}", collectorKey, minDoc, maxDoc, bits.cardinality(), - wordCount + wordCount, + nextDoc ); } - return wordCount; + return ((long) nextDoc << 32) | (wordCount & 0xFFFFFFFFL); } @Override From 42a1e959d4ce625936f703251e0c6c7eb0e2f58d Mon Sep 17 00:00:00 2001 From: Andriy Redko Date: Thu, 20 Aug 2026 07:39:12 -0400 Subject: [PATCH 5/6] Remove Jackson 2.x dependencies from OpenSearch core (#22704) * Remove Jackson 2.x dependencies from OpenSearch core Signed-off-by: Andriy Redko * Address review comments Signed-off-by: Andriy Redko --------- Signed-off-by: Andriy Redko --- client/sniffer/licenses/jackson-core-LICENSE | 2 +- libs/core/licenses/jackson-LICENSE | 2 +- libs/x-content/licenses/jackson-LICENSE | 2 +- modules/ingest-geoip/build.gradle | 6 ++++++ ...on-annotations-LICENSE => jackson-LICENSE} | 2 +- ...kson-annotations-NOTICE => jackson-NOTICE} | 0 .../licenses/jackson-core-2.22.1.jar.sha1 | 0 plugins/arrow-base/build.gradle | 1 + plugins/arrow-base/licenses/jackson-LICENSE | 2 +- .../licenses/jackson-core-2.22.1.jar.sha1 | 1 + plugins/crypto-kms/licenses/jackson-LICENSE | 2 +- plugins/discovery-ec2/build.gradle | 1 + .../discovery-ec2/licenses/jackson-LICENSE | 2 +- .../licenses/jackson-core-2.22.1.jar.sha1 | 1 + plugins/discovery-gce/build.gradle | 3 +++ .../discovery-gce/licenses/jackson-LICENSE | 2 +- .../discovery-gce/licenses/jackson-NOTICE | 0 .../jackson-annotations-2.22.jar.sha1 | 1 + .../licenses/jackson-core-2.22.1.jar.sha1 | 1 + .../licenses/jackson-LICENSE | 2 +- .../licenses/jackson-core-2.22.1.jar.sha1 | 1 + .../jackson-dataformat-cbor-2.22.1.jar.sha1 | 0 plugins/repository-azure/build.gradle | 1 + .../licenses/jackson-core-2.22.1.jar.sha1 | 1 + plugins/repository-gcs/build.gradle | 2 ++ .../repository-gcs/licenses/jackson-LICENSE | 2 +- .../repository-gcs/licenses/jackson-NOTICE | 9 ++++++--- .../jackson-annotations-2.22.jar.sha1 | 1 + .../licenses/jackson-core-2.22.1.jar.sha1 | 1 + .../repository-s3/licenses/jackson-LICENSE | 2 +- .../licenses/jackson-core-2.22.1.jar.sha1 | 1 + .../jackson-dataformat-cbor-2.22.1.jar.sha1 | 1 + qa/wildfly/build.gradle | 1 + server/build.gradle | 20 +++++++++++-------- server/licenses/jackson-LICENSE | 2 +- .../jackson-dataformat-smile-2.22.1.jar.sha1 | 1 - .../jackson-dataformat-yaml-2.22.1.jar.sha1 | 1 - server/licenses/snakeyaml-2.6.jar.sha1 | 1 - 38 files changed, 55 insertions(+), 26 deletions(-) rename modules/ingest-geoip/licenses/{jackson-annotations-LICENSE => jackson-LICENSE} (86%) rename modules/ingest-geoip/licenses/{jackson-annotations-NOTICE => jackson-NOTICE} (100%) rename {server => modules/ingest-geoip}/licenses/jackson-core-2.22.1.jar.sha1 (100%) create mode 100644 plugins/arrow-base/licenses/jackson-core-2.22.1.jar.sha1 create mode 100644 plugins/discovery-ec2/licenses/jackson-core-2.22.1.jar.sha1 rename modules/ingest-geoip/licenses/jackson-databind-LICENSE => plugins/discovery-gce/licenses/jackson-LICENSE (86%) rename modules/ingest-geoip/licenses/jackson-databind-NOTICE => plugins/discovery-gce/licenses/jackson-NOTICE (100%) create mode 100644 plugins/discovery-gce/licenses/jackson-annotations-2.22.jar.sha1 create mode 100644 plugins/discovery-gce/licenses/jackson-core-2.22.1.jar.sha1 create mode 100644 plugins/ingestion-kinesis/licenses/jackson-core-2.22.1.jar.sha1 rename {server => plugins/ingestion-kinesis}/licenses/jackson-dataformat-cbor-2.22.1.jar.sha1 (100%) create mode 100644 plugins/repository-azure/licenses/jackson-core-2.22.1.jar.sha1 rename modules/ingest-geoip/licenses/jackson-datatype-jsr310-LICENSE => plugins/repository-gcs/licenses/jackson-LICENSE (74%) rename modules/ingest-geoip/licenses/jackson-datatype-jsr310-NOTICE => plugins/repository-gcs/licenses/jackson-NOTICE (53%) create mode 100644 plugins/repository-gcs/licenses/jackson-annotations-2.22.jar.sha1 create mode 100644 plugins/repository-gcs/licenses/jackson-core-2.22.1.jar.sha1 create mode 100644 plugins/repository-s3/licenses/jackson-core-2.22.1.jar.sha1 create mode 100644 plugins/repository-s3/licenses/jackson-dataformat-cbor-2.22.1.jar.sha1 delete mode 100644 server/licenses/jackson-dataformat-smile-2.22.1.jar.sha1 delete mode 100644 server/licenses/jackson-dataformat-yaml-2.22.1.jar.sha1 delete mode 100644 server/licenses/snakeyaml-2.6.jar.sha1 diff --git a/client/sniffer/licenses/jackson-core-LICENSE b/client/sniffer/licenses/jackson-core-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/client/sniffer/licenses/jackson-core-LICENSE +++ b/client/sniffer/licenses/jackson-core-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/libs/core/licenses/jackson-LICENSE b/libs/core/licenses/jackson-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/libs/core/licenses/jackson-LICENSE +++ b/libs/core/licenses/jackson-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/libs/x-content/licenses/jackson-LICENSE b/libs/x-content/licenses/jackson-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/libs/x-content/licenses/jackson-LICENSE +++ b/libs/x-content/licenses/jackson-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/modules/ingest-geoip/build.gradle b/modules/ingest-geoip/build.gradle index 99c757989c6d6..5bd4cc40de374 100644 --- a/modules/ingest-geoip/build.gradle +++ b/modules/ingest-geoip/build.gradle @@ -42,7 +42,9 @@ dependencies { api('com.maxmind.geoip2:geoip2:4.4.0') // geoip2 dependencies: api('com.maxmind.db:maxmind-db:3.2.0') + // required by geoip2 api(libs.jackson.annotation) + api(libs.jackson.core) api(libs.jackson.databind) api(libs.jackson.datatype.jsr310) @@ -75,3 +77,7 @@ if (Os.isFamily(Os.FAMILY_WINDOWS)) { systemProperty 'opensearch.geoip.load_db_on_heap', 'true' } } + +tasks.named("dependencyLicenses").configure { + mapping from: /jackson-.*/, to: 'jackson' +} diff --git a/modules/ingest-geoip/licenses/jackson-annotations-LICENSE b/modules/ingest-geoip/licenses/jackson-LICENSE similarity index 86% rename from modules/ingest-geoip/licenses/jackson-annotations-LICENSE rename to modules/ingest-geoip/licenses/jackson-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/modules/ingest-geoip/licenses/jackson-annotations-LICENSE +++ b/modules/ingest-geoip/licenses/jackson-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/modules/ingest-geoip/licenses/jackson-annotations-NOTICE b/modules/ingest-geoip/licenses/jackson-NOTICE similarity index 100% rename from modules/ingest-geoip/licenses/jackson-annotations-NOTICE rename to modules/ingest-geoip/licenses/jackson-NOTICE diff --git a/server/licenses/jackson-core-2.22.1.jar.sha1 b/modules/ingest-geoip/licenses/jackson-core-2.22.1.jar.sha1 similarity index 100% rename from server/licenses/jackson-core-2.22.1.jar.sha1 rename to modules/ingest-geoip/licenses/jackson-core-2.22.1.jar.sha1 diff --git a/plugins/arrow-base/build.gradle b/plugins/arrow-base/build.gradle index 8ea3135740b5c..2cba5edf8a7f0 100644 --- a/plugins/arrow-base/build.gradle +++ b/plugins/arrow-base/build.gradle @@ -30,6 +30,7 @@ dependencies { api "com.google.flatbuffers:flatbuffers-java:${versions.flatbuffers}" api "org.slf4j:slf4j-api:${versions.slf4j}" api "com.fasterxml.jackson.core:jackson-annotations:${versions.jackson_annotations}" + api "com.fasterxml.jackson.core:jackson-core:${versions.jackson}" api "com.fasterxml.jackson.core:jackson-databind:${versions.jackson}" api "com.fasterxml.jackson.datatype:jackson-datatype-jsr310:${versions.jackson}" api "commons-codec:commons-codec:${versions.commonscodec}" diff --git a/plugins/arrow-base/licenses/jackson-LICENSE b/plugins/arrow-base/licenses/jackson-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/plugins/arrow-base/licenses/jackson-LICENSE +++ b/plugins/arrow-base/licenses/jackson-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/plugins/arrow-base/licenses/jackson-core-2.22.1.jar.sha1 b/plugins/arrow-base/licenses/jackson-core-2.22.1.jar.sha1 new file mode 100644 index 0000000000000..e2c41405d169c --- /dev/null +++ b/plugins/arrow-base/licenses/jackson-core-2.22.1.jar.sha1 @@ -0,0 +1 @@ +da7ffb60088d7e8f37ecdd3b617520971cc7b9bf \ No newline at end of file diff --git a/plugins/crypto-kms/licenses/jackson-LICENSE b/plugins/crypto-kms/licenses/jackson-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/plugins/crypto-kms/licenses/jackson-LICENSE +++ b/plugins/crypto-kms/licenses/jackson-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/plugins/discovery-ec2/build.gradle b/plugins/discovery-ec2/build.gradle index db459a53f556e..8afb3b4e981d0 100644 --- a/plugins/discovery-ec2/build.gradle +++ b/plugins/discovery-ec2/build.gradle @@ -73,6 +73,7 @@ dependencies { api "org.apache.logging.log4j:log4j-1.2-api:${versions.log4j}" api "org.slf4j:slf4j-api:${versions.slf4j}" api "commons-codec:commons-codec:${versions.commonscodec}" + api "com.fasterxml.jackson.core:jackson-core:${versions.jackson}" api "com.fasterxml.jackson.core:jackson-databind:${versions.jackson_databind}" api "com.fasterxml.jackson.core:jackson-annotations:${versions.jackson_annotations}" api "org.reactivestreams:reactive-streams:${versions.reactivestreams}" diff --git a/plugins/discovery-ec2/licenses/jackson-LICENSE b/plugins/discovery-ec2/licenses/jackson-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/plugins/discovery-ec2/licenses/jackson-LICENSE +++ b/plugins/discovery-ec2/licenses/jackson-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/plugins/discovery-ec2/licenses/jackson-core-2.22.1.jar.sha1 b/plugins/discovery-ec2/licenses/jackson-core-2.22.1.jar.sha1 new file mode 100644 index 0000000000000..e2c41405d169c --- /dev/null +++ b/plugins/discovery-ec2/licenses/jackson-core-2.22.1.jar.sha1 @@ -0,0 +1 @@ +da7ffb60088d7e8f37ecdd3b617520971cc7b9bf \ No newline at end of file diff --git a/plugins/discovery-gce/build.gradle b/plugins/discovery-gce/build.gradle index b62a6e0af77a9..211e5320faf03 100644 --- a/plugins/discovery-gce/build.gradle +++ b/plugins/discovery-gce/build.gradle @@ -25,6 +25,8 @@ dependencies { api "com.google.http-client:google-http-client:${versions.google_http_client}" api "com.google.http-client:google-http-client-gson:${versions.google_http_client}" api "com.google.http-client:google-http-client-jackson2:${versions.google_http_client}" + api "com.fasterxml.jackson.core:jackson-annotations:${versions.jackson_annotations}" + api "com.fasterxml.jackson.core:jackson-core:${versions.jackson}" api 'com.google.code.findbugs:jsr305:3.0.2' api "org.apache.httpcomponents:httpclient:${versions.httpclient}" api "org.apache.httpcomponents:httpcore:${versions.httpcore}" @@ -49,6 +51,7 @@ restResources { tasks.named("dependencyLicenses").configure { mapping from: /google-.*/, to: 'google' mapping from: /opencensus.*/, to: 'opencensus' + mapping from: /jackson-.*/, to: 'jackson' } check { diff --git a/modules/ingest-geoip/licenses/jackson-databind-LICENSE b/plugins/discovery-gce/licenses/jackson-LICENSE similarity index 86% rename from modules/ingest-geoip/licenses/jackson-databind-LICENSE rename to plugins/discovery-gce/licenses/jackson-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/modules/ingest-geoip/licenses/jackson-databind-LICENSE +++ b/plugins/discovery-gce/licenses/jackson-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/modules/ingest-geoip/licenses/jackson-databind-NOTICE b/plugins/discovery-gce/licenses/jackson-NOTICE similarity index 100% rename from modules/ingest-geoip/licenses/jackson-databind-NOTICE rename to plugins/discovery-gce/licenses/jackson-NOTICE diff --git a/plugins/discovery-gce/licenses/jackson-annotations-2.22.jar.sha1 b/plugins/discovery-gce/licenses/jackson-annotations-2.22.jar.sha1 new file mode 100644 index 0000000000000..de60cbfe0677d --- /dev/null +++ b/plugins/discovery-gce/licenses/jackson-annotations-2.22.jar.sha1 @@ -0,0 +1 @@ +15c67f9498d6934cff9510fc70c5a12e52290457 \ No newline at end of file diff --git a/plugins/discovery-gce/licenses/jackson-core-2.22.1.jar.sha1 b/plugins/discovery-gce/licenses/jackson-core-2.22.1.jar.sha1 new file mode 100644 index 0000000000000..e2c41405d169c --- /dev/null +++ b/plugins/discovery-gce/licenses/jackson-core-2.22.1.jar.sha1 @@ -0,0 +1 @@ +da7ffb60088d7e8f37ecdd3b617520971cc7b9bf \ No newline at end of file diff --git a/plugins/ingestion-kinesis/licenses/jackson-LICENSE b/plugins/ingestion-kinesis/licenses/jackson-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/plugins/ingestion-kinesis/licenses/jackson-LICENSE +++ b/plugins/ingestion-kinesis/licenses/jackson-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/plugins/ingestion-kinesis/licenses/jackson-core-2.22.1.jar.sha1 b/plugins/ingestion-kinesis/licenses/jackson-core-2.22.1.jar.sha1 new file mode 100644 index 0000000000000..e2c41405d169c --- /dev/null +++ b/plugins/ingestion-kinesis/licenses/jackson-core-2.22.1.jar.sha1 @@ -0,0 +1 @@ +da7ffb60088d7e8f37ecdd3b617520971cc7b9bf \ No newline at end of file diff --git a/server/licenses/jackson-dataformat-cbor-2.22.1.jar.sha1 b/plugins/ingestion-kinesis/licenses/jackson-dataformat-cbor-2.22.1.jar.sha1 similarity index 100% rename from server/licenses/jackson-dataformat-cbor-2.22.1.jar.sha1 rename to plugins/ingestion-kinesis/licenses/jackson-dataformat-cbor-2.22.1.jar.sha1 diff --git a/plugins/repository-azure/build.gradle b/plugins/repository-azure/build.gradle index 8595a26851a02..5651109fb92fb 100644 --- a/plugins/repository-azure/build.gradle +++ b/plugins/repository-azure/build.gradle @@ -77,6 +77,7 @@ dependencies { api "io.projectreactor.netty:reactor-netty-http:${versions.reactor_netty}" api "org.slf4j:slf4j-api:${versions.slf4j}" api "com.fasterxml.jackson.core:jackson-annotations:${versions.jackson_annotations}" + api "com.fasterxml.jackson.core:jackson-core:${versions.jackson}" api "com.fasterxml.jackson.core:jackson-databind:${versions.jackson_databind}" api "com.fasterxml.jackson.datatype:jackson-datatype-jsr310:${versions.jackson}" api "com.fasterxml.jackson.dataformat:jackson-dataformat-xml:${versions.jackson}" diff --git a/plugins/repository-azure/licenses/jackson-core-2.22.1.jar.sha1 b/plugins/repository-azure/licenses/jackson-core-2.22.1.jar.sha1 new file mode 100644 index 0000000000000..e2c41405d169c --- /dev/null +++ b/plugins/repository-azure/licenses/jackson-core-2.22.1.jar.sha1 @@ -0,0 +1 @@ +da7ffb60088d7e8f37ecdd3b617520971cc7b9bf \ No newline at end of file diff --git a/plugins/repository-gcs/build.gradle b/plugins/repository-gcs/build.gradle index 1552852532c3a..98d453763c62b 100644 --- a/plugins/repository-gcs/build.gradle +++ b/plugins/repository-gcs/build.gradle @@ -83,6 +83,7 @@ dependencies { runtimeOnly "com.google.http-client:google-http-client-gson:2.0.3" runtimeOnly "com.google.http-client:google-http-client-appengine:2.0.3" runtimeOnly "com.google.http-client:google-http-client-jackson2:2.0.3" + runtimeOnly "com.fasterxml.jackson.core:jackson-annotations:${versions.jackson_annotations}" runtimeOnly "com.fasterxml.jackson.core:jackson-core:${versions.jackson}" // 2.18.2 in bom runtimeOnly "io.opencensus:opencensus-api:0.31.1" runtimeOnly "io.opencensus:opencensus-contrib-http-util:0.31.1" @@ -141,6 +142,7 @@ tasks.named("dependencyLicenses").configure { mapping from: /opentelemetry-common.*/, to: 'opentelemetry-api' mapping from: /protobuf.*/, to: 'protobuf' mapping from: /proto-google.*/, to: 'proto-google' + mapping from: /jackson-.*/, to: 'jackson' } thirdPartyAudit { diff --git a/modules/ingest-geoip/licenses/jackson-datatype-jsr310-LICENSE b/plugins/repository-gcs/licenses/jackson-LICENSE similarity index 74% rename from modules/ingest-geoip/licenses/jackson-datatype-jsr310-LICENSE rename to plugins/repository-gcs/licenses/jackson-LICENSE index 0e9c9520aeeee..893b236ee25f2 100644 --- a/modules/ingest-geoip/licenses/jackson-datatype-jsr310-LICENSE +++ b/plugins/repository-gcs/licenses/jackson-LICENSE @@ -1,4 +1,4 @@ -This copy of Jackson JSON processor Java 8 Date/Time module is licensed under the +This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the specific rights regarding derivative works. diff --git a/modules/ingest-geoip/licenses/jackson-datatype-jsr310-NOTICE b/plugins/repository-gcs/licenses/jackson-NOTICE similarity index 53% rename from modules/ingest-geoip/licenses/jackson-datatype-jsr310-NOTICE rename to plugins/repository-gcs/licenses/jackson-NOTICE index d55c59a0d506f..4c976b7b4cc58 100644 --- a/modules/ingest-geoip/licenses/jackson-datatype-jsr310-NOTICE +++ b/plugins/repository-gcs/licenses/jackson-NOTICE @@ -3,12 +3,15 @@ Jackson is a high-performance, Free/Open Source JSON processing library. It was originally written by Tatu Saloranta (tatu.saloranta@iki.fi), and has been in development since 2007. -It is currently developed by a community of developers. +It is currently developed by a community of developers, as well as supported +commercially by FasterXML.com. ## Licensing -Jackson components are licensed under Apache (Software) License, version 2.0, -as per accompanying LICENSE file. +Jackson core and extension components may licensed under different licenses. +To find the details that apply to this artifact see the accompanying LICENSE file. +For more information, including possible other licensing options, contact +FasterXML.com (http://fasterxml.com). ## Credits diff --git a/plugins/repository-gcs/licenses/jackson-annotations-2.22.jar.sha1 b/plugins/repository-gcs/licenses/jackson-annotations-2.22.jar.sha1 new file mode 100644 index 0000000000000..de60cbfe0677d --- /dev/null +++ b/plugins/repository-gcs/licenses/jackson-annotations-2.22.jar.sha1 @@ -0,0 +1 @@ +15c67f9498d6934cff9510fc70c5a12e52290457 \ No newline at end of file diff --git a/plugins/repository-gcs/licenses/jackson-core-2.22.1.jar.sha1 b/plugins/repository-gcs/licenses/jackson-core-2.22.1.jar.sha1 new file mode 100644 index 0000000000000..e2c41405d169c --- /dev/null +++ b/plugins/repository-gcs/licenses/jackson-core-2.22.1.jar.sha1 @@ -0,0 +1 @@ +da7ffb60088d7e8f37ecdd3b617520971cc7b9bf \ No newline at end of file diff --git a/plugins/repository-s3/licenses/jackson-LICENSE b/plugins/repository-s3/licenses/jackson-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/plugins/repository-s3/licenses/jackson-LICENSE +++ b/plugins/repository-s3/licenses/jackson-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/plugins/repository-s3/licenses/jackson-core-2.22.1.jar.sha1 b/plugins/repository-s3/licenses/jackson-core-2.22.1.jar.sha1 new file mode 100644 index 0000000000000..e2c41405d169c --- /dev/null +++ b/plugins/repository-s3/licenses/jackson-core-2.22.1.jar.sha1 @@ -0,0 +1 @@ +da7ffb60088d7e8f37ecdd3b617520971cc7b9bf \ No newline at end of file diff --git a/plugins/repository-s3/licenses/jackson-dataformat-cbor-2.22.1.jar.sha1 b/plugins/repository-s3/licenses/jackson-dataformat-cbor-2.22.1.jar.sha1 new file mode 100644 index 0000000000000..a6e9a8c626cfe --- /dev/null +++ b/plugins/repository-s3/licenses/jackson-dataformat-cbor-2.22.1.jar.sha1 @@ -0,0 +1 @@ +f1a9be02f8e042ee37441116b31104cec8ec9373 \ No newline at end of file diff --git a/qa/wildfly/build.gradle b/qa/wildfly/build.gradle index 12f2eacccd982..101cc2aa0c154 100644 --- a/qa/wildfly/build.gradle +++ b/qa/wildfly/build.gradle @@ -54,6 +54,7 @@ dependencies { exclude group: 'com.fasterxml.jackson.dataformat' exclude group: 'com.fasterxml.jackson.module' } + api "com.fasterxml.jackson.core:jackson-core:${versions.jackson}" api "com.fasterxml.jackson.core:jackson-annotations:${versions.jackson_annotations}" api "com.fasterxml.jackson.core:jackson-databind:${versions.jackson_databind}" api "com.fasterxml.jackson.jakarta.rs:jackson-jakarta-rs-base:${versions.jackson}" diff --git a/server/build.gradle b/server/build.gradle index ad38c8b21ed02..2b47f60cd31f5 100644 --- a/server/build.gradle +++ b/server/build.gradle @@ -122,13 +122,6 @@ dependencies { api libs.protobuf api libs.jakartaannotation - // Keeping Jackson 2.x dependencies till migration to 3.x completes across the ecosystem - api libs.jackson.core - api libs.jackson.dataformat.smile - api libs.jackson.dataformat.yaml - api libs.jackson.dataformat.cbor - api "org.yaml:snakeyaml:${versions.snakeyaml}" - // https://mvnrepository.com/artifact/org.roaringbitmap/RoaringBitmap api libs.roaringbitmap testImplementation 'org.awaitility:awaitility:4.3.0' @@ -256,10 +249,21 @@ tasks.named("thirdPartyAudit").configure { // from com.fasterxml.jackson.dataformat.yaml.YAMLMapper (jackson-dataformat-yaml) 'com.fasterxml.jackson.databind.ObjectMapper', 'com.fasterxml.jackson.databind.util.ClassUtil', - 'com.fasterxml.jackson.databind.cfg.MapperBuilder', // from log4j 'com.conversantmedia.util.concurrent.SpinPolicy', + 'com.fasterxml.jackson.core.JsonGenerator', + 'com.fasterxml.jackson.core.JsonParser', + 'com.fasterxml.jackson.core.JsonParser$Feature', + 'com.fasterxml.jackson.core.JsonToken', + 'com.fasterxml.jackson.core.ObjectCodec', + 'com.fasterxml.jackson.core.PrettyPrinter', + 'com.fasterxml.jackson.core.Version', + 'com.fasterxml.jackson.core.Versioned', + 'com.fasterxml.jackson.core.io.IOContext', + 'com.fasterxml.jackson.core.type.TypeReference', + 'com.fasterxml.jackson.core.util.VersionUtil', + 'com.fasterxml.jackson.dataformat.yaml.YAMLMapper', 'com.fasterxml.jackson.annotation.JsonInclude$Include', 'com.fasterxml.jackson.databind.DeserializationContext', 'com.fasterxml.jackson.databind.DeserializationFeature', diff --git a/server/licenses/jackson-LICENSE b/server/licenses/jackson-LICENSE index f5f45d26a49d6..893b236ee25f2 100644 --- a/server/licenses/jackson-LICENSE +++ b/server/licenses/jackson-LICENSE @@ -1,7 +1,7 @@ This copy of Jackson JSON processor streaming parser/generator is licensed under the Apache (Software) License, version 2.0 ("the License"). See the License for details about distribution rights, and the -specific rights regarding derivate works. +specific rights regarding derivative works. You may obtain a copy of the License at: diff --git a/server/licenses/jackson-dataformat-smile-2.22.1.jar.sha1 b/server/licenses/jackson-dataformat-smile-2.22.1.jar.sha1 deleted file mode 100644 index 7986a3d566008..0000000000000 --- a/server/licenses/jackson-dataformat-smile-2.22.1.jar.sha1 +++ /dev/null @@ -1 +0,0 @@ -82b68722e819e163f15895219dc748f88db95a6c \ No newline at end of file diff --git a/server/licenses/jackson-dataformat-yaml-2.22.1.jar.sha1 b/server/licenses/jackson-dataformat-yaml-2.22.1.jar.sha1 deleted file mode 100644 index 597173dcd7fc3..0000000000000 --- a/server/licenses/jackson-dataformat-yaml-2.22.1.jar.sha1 +++ /dev/null @@ -1 +0,0 @@ -af5fde2414e4a8d2617e7890f7068c26b5b66a33 \ No newline at end of file diff --git a/server/licenses/snakeyaml-2.6.jar.sha1 b/server/licenses/snakeyaml-2.6.jar.sha1 deleted file mode 100644 index 65f4b2277980b..0000000000000 --- a/server/licenses/snakeyaml-2.6.jar.sha1 +++ /dev/null @@ -1 +0,0 @@ -2bc14918a2f8d5414749ab12d0c590cd3198b8c1 \ No newline at end of file From f3b7c277895fad542012d65d50f8b6fec1c60118 Mon Sep 17 00:00:00 2001 From: Darshit Chanpura Date: Fri, 21 Aug 2026 02:58:10 +0000 Subject: [PATCH 6/6] Add unit test for constant_score preservation in buildFilteredQuery Covers both branches of the ConstantScoreQuery re-wrap in DefaultSearchContext.buildFilteredQuery: a ConstantScoreQuery input with a filter present stays non-scoring (returns a ConstantScoreQuery), and a normal query is combined into a scoring BooleanQuery. Complements the existing ConstantScoreFilteredAliasIT with method-level coverage. Signed-off-by: Darshit Chanpura --- .../search/DefaultSearchContextTests.java | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/server/src/test/java/org/opensearch/search/DefaultSearchContextTests.java b/server/src/test/java/org/opensearch/search/DefaultSearchContextTests.java index b6c236654569a..1cb338483de14 100644 --- a/server/src/test/java/org/opensearch/search/DefaultSearchContextTests.java +++ b/server/src/test/java/org/opensearch/search/DefaultSearchContextTests.java @@ -35,12 +35,16 @@ import com.carrotsearch.randomizedtesting.annotations.ParametersFactory; import org.apache.lucene.index.IndexReader; +import org.apache.lucene.index.Term; +import org.apache.lucene.search.BooleanQuery; +import org.apache.lucene.search.ConstantScoreQuery; import org.apache.lucene.search.IndexSearcher; import org.apache.lucene.search.MatchNoDocsQuery; import org.apache.lucene.search.Query; import org.apache.lucene.search.QueryCachingPolicy; import org.apache.lucene.search.Sort; import org.apache.lucene.search.SortField; +import org.apache.lucene.search.TermQuery; import org.apache.lucene.store.Directory; import org.apache.lucene.tests.index.RandomIndexWriter; import org.opensearch.Version; @@ -111,6 +115,7 @@ import static org.opensearch.index.IndexSettings.INDEX_SEARCH_THROTTLED; import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.instanceOf; import static org.hamcrest.Matchers.is; import static org.mockito.Mockito.any; import static org.mockito.Mockito.anyBoolean; @@ -398,6 +403,20 @@ protected Engine.Searcher acquireSearcherInternal(String source) { ParsedQuery parsedQuery = ParsedQuery.parsedMatchAllQuery(); context3.sliceBuilder(null).parsedQuery(parsedQuery).preProcess(false); assertEquals(context3.query(), context3.buildFilteredQuery(parsedQuery.query())); + + // buildFilteredQuery: when at least one filter is added (here a slice filter) and the incoming + // query is already a ConstantScoreQuery, the combined query must stay non-scoring so Lucene's + // COMPLETE_NO_SCORES fast path is preserved; a normal query is combined into a scoring BooleanQuery. + when(mapperService.hasNested()).thenReturn(false); + SliceBuilder filterSlice = mock(SliceBuilder.class); + when(filterSlice.toFilter(any(), any(), any(), any())).thenReturn(new TermQuery(new Term("field", "value"))); + context3.sliceBuilder(filterSlice); + assertThat( + context3.buildFilteredQuery(new ConstantScoreQuery(new TermQuery(new Term("content", "alpha")))), + instanceOf(ConstantScoreQuery.class) + ); + assertThat(context3.buildFilteredQuery(new TermQuery(new Term("content", "alpha"))), instanceOf(BooleanQuery.class)); + context3.sliceBuilder(null); // make sure getPreciseRelativeTimeInMillis is same as System.nanoTime() long timeToleranceInMs = 10; long currTime = TimeValue.nsecToMSec(System.nanoTime());