diff --git a/.gitignore b/.gitignore index f7ee48d3..7f280706 100644 --- a/.gitignore +++ b/.gitignore @@ -7,3 +7,5 @@ target bin .project cuvs-workdir +__pycache__/ +.pytest_cache/ diff --git a/README.md b/README.md index 4162e19d..472eb23b 100644 --- a/README.md +++ b/README.md @@ -104,6 +104,147 @@ public class HelloCuvsLucene { The artifacts would be built and available in the target / folder. +### Using with PyLucene + +PyLucene embeds a JVM and starts it with the classpath passed to `lucene.initVM(...)`. +Because PyLucene's generated Python module only exposes the Java classes it was built +to wrap, use Lucene's service provider lookup to load `cuvs-lucene` codecs from +Python instead of importing `com.nvidia.cuvs.lucene` classes directly. + +Build the standard cuvs-lucene jar: + +```sh +mvn clean package -DskipTests +``` + +Then start PyLucene with the base `cuvs-java` jar, the standard `cuvs-lucene` +jar, and PyLucene's own Lucene classpath: + +```python +import os +from pathlib import Path + +import lucene + +cuvs_java_jar = Path(os.environ["CUVS_LUCENE_CUVS_JAVA_JAR"]) +cuvs_lucene_jar = next( + jar + for jar in Path("target").glob("cuvs-lucene-*.jar") + if "-jar-with-" not in jar.name + and not jar.name.endswith(("-sources.jar", "-javadoc.jar")) +) +lucene.initVM( + classpath=os.pathsep.join( + [str(cuvs_java_jar), str(cuvs_lucene_jar), lucene.CLASSPATH] + ), + vmargs=[ + "--enable-native-access=ALL-UNNAMED", + "--add-modules=jdk.incubator.vector", + ], +) + +from org.apache.lucene.codecs import Codec + +codec = Codec.forName("Lucene101AcceleratedHNSWCodec") +``` + +Use the returned `codec` with `IndexWriterConfig.setCodec(codec)`. The standard +artifact includes `cuvs-lucene` classes and service descriptors. +PyLucene must provide Lucene classes, and the base multi-release `cuvs-java` jar +must be present separately on the JVM classpath. Do not use a native classifier +`cuvs-java` jar here unless you also want to rely on its embedded native +libraries; the base jar uses native libraries from +`LD_LIBRARY_PATH`/`java.library.path`. + +#### PyLucene end-to-end tests + +Pytest cases are under `src/test/python`. The parametrized cases and assertions +are in `test_pylucene_end_to_end.py`; reusable index and search code is in +`pylucene_test_support.py`. The test-only Java adapters in +`src/test/java/com/nvidia/cuvs/lucene/PyLuceneTestSupport.java` are compiled to +`target/test-classes` and are not included in the published jar. + +`ci/run_pylucene_pytests.sh` builds the Maven artifacts when requested, resolves +the PyLucene classpath inputs, and invokes pytest. `test_pylucene.sh` at the +repository root is the convenient entry point. + +To build the artifacts and run the default CPU check in an activated PyLucene +environment: + +```sh +./test_pylucene.sh +``` + +The default runs the jar-packaging checks and `cpu-hnsw-1-segment`. Use the +groups below for broader coverage. +Use `--no-build` when the standard jar and test bridge are already compiled. +`CUVS_LUCENE_PYLUCENE_TEST_CLASSES` can point to a different test-classes +directory and defaults to `target/test-classes`. + +Pytest can also be invoked directly once the PyLucene environment and classpath +inputs are available: + +```sh +CUVS_LUCENE_JAR=/absolute/path/to/cuvs-lucene.jar \ +CUVS_LUCENE_CUVS_JAVA_JAR=/absolute/path/to/cuvs-java.jar \ +CUVS_LUCENE_PYLUCENE_TEST_CLASSES="$(pwd)/target/test-classes" \ +python3 -m pytest -q -s src/test/python/test_pylucene_end_to_end.py +``` + +The execution-path groups are: + +- `cpu-hnsw`: HNSW build and search through Lucene's CPU path. +- `gpu-cagra-built-hnsw`: GPU CAGRA build followed by HNSW search. +- `gpu-cagra-search`: GPU CAGRA build and search. + +GPU cases assert that cuVS was actually used and fail if it is unavailable or +falls back to CPU. CPU cases assert and report the CPU path, and HNSW cases +verify the persisted graph shape. CAGRA cases use `graphDegree=32` and +`intermediateGraphDegree=64`. + +Vectors and queries are deterministic. Expected neighbors are computed with +brute force. Queries for live indexed vectors check rank-one self matches, +duplicate hits, and a configurable recall floor. Separate tests cover segment +topology, force merges, HNSW layer count, CAGRA `searchWidth` values 1, 16, and +32, and document filters. A dedicated CAGRA-search case verifies that a deleted +document is not returned. Set the recall floor with `--min-recall=FLOAT` in the +wrapper or `CUVS_LUCENE_PYLUCENE_MIN_RECALL` for direct pytest execution. +The default floor is `0.75`. + +The `document-filter` group exercises selective filters through CPU HNSW, +CAGRA-built HNSW, and CAGRA search. The CAGRA-search case uses ten segments +and accepts roughly one quarter of each segment. Every segment retains more +than `topK` accepted vectors so Lucene exercises approximate native +prefiltering; results are checked against brute-force neighbors from only the +accepted vectors. + +Useful behavior groups include `execution-paths`, `segment-topologies`, +`force-merges`, `hnsw-layer-counts`, `cagra-search-widths`, +`deleted-documents`, and `document-filter`. + +Run the complete CPU/GPU end-to-end suite with: + +```sh +./test_pylucene.sh --full-e2e +``` + +Select a focused group or set the minimum document count per scenario: + +```sh +./test_pylucene.sh --cases=cagra-search-widths \ + --rows=5000 --dims=64 --topk=20 --min-recall=0.8 +``` + +`--rows` is a lower bound. CAGRA cases use at least 97 vector-bearing +documents per constructed graph, one more than cuVS's internal NN-Descent +degree of 96. The three-layer HNSW case uses at least 24,832 vector-bearing +documents so its third layer retains 97. Arguments after `--` are passed to +pytest, for example: + +```sh +./test_pylucene.sh --no-build --cases=gpu-cagra-search -- -x +``` + ### Running Tests ```sh diff --git a/ci/run_pylucene_pytests.sh b/ci/run_pylucene_pytests.sh new file mode 100755 index 00000000..f222f353 --- /dev/null +++ b/ci/run_pylucene_pytests.sh @@ -0,0 +1,232 @@ +#!/bin/bash + +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +set -euo pipefail + +REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +cd "${REPO_ROOT}" + +MVN_BIN="${MVN:-mvn}" +PYTHON_BIN="${PYTHON:-python3}" +SKIP_BUILD=0 +FULL_END_TO_END=0 +PYLUCENE_CASES="${CUVS_LUCENE_PYLUCENE_CASES:-}" +PYLUCENE_ROWS="${CUVS_LUCENE_PYLUCENE_ROWS:-}" +PYLUCENE_DIMS="${CUVS_LUCENE_PYLUCENE_DIMS:-}" +PYLUCENE_TOPK="${CUVS_LUCENE_PYLUCENE_TOPK:-}" +PYLUCENE_MIN_RECALL="${CUVS_LUCENE_PYLUCENE_MIN_RECALL:-}" +CUVS_LUCENE_JAR_PATH="${CUVS_LUCENE_JAR:-}" +PYLUCENE_TEST_CLASSES_PATH="${CUVS_LUCENE_PYLUCENE_TEST_CLASSES:-target/test-classes}" +PYTEST_ARGS=() + +while [[ "$#" -gt 0 ]]; do + arg="$1" + shift + case "${arg}" in + --full-e2e|--gpu-e2e) + FULL_END_TO_END=1 + ;; + --no-build) + SKIP_BUILD=1 + ;; + --cases=*) + PYLUCENE_CASES="${arg#--cases=}" + ;; + --rows=*) + PYLUCENE_ROWS="${arg#--rows=}" + ;; + --dims=*) + PYLUCENE_DIMS="${arg#--dims=}" + ;; + --topk=*) + PYLUCENE_TOPK="${arg#--topk=}" + ;; + --min-recall=*) + PYLUCENE_MIN_RECALL="${arg#--min-recall=}" + ;; + --) + PYTEST_ARGS+=("$@") + break + ;; + -h|--help) + echo "Usage: ./test_pylucene.sh [OPTIONS] [-- PYTEST_ARGS...]" + echo + echo "Prepares the PyLucene JVM, jar, classpath, and native-library inputs, then runs pytest." + echo + echo "Options:" + echo " --no-build Use existing Maven artifacts." + echo " --full-e2e Run the complete CPU/GPU end-to-end suite." + echo " --cases=CASE[,CASE...] Select case names or groups." + echo " --rows=N Set the minimum document count per scenario." + echo " --dims=N Set vector dimensions." + echo " --topk=N Set requested neighbor count." + echo " --min-recall=FLOAT Set the brute-force recall floor (default: 0.75)." + echo " -- PYTEST_ARGS Forward remaining arguments to pytest." + echo + echo "Execution-path groups:" + echo " cpu-hnsw, gpu-cagra-built-hnsw, gpu-cagra-search, gpu" + echo + echo "Behavior groups:" + echo " execution-paths, segment-topologies, force-merges," + echo " hnsw-layer-counts, cagra-search-widths," + echo " deleted-documents, document-filter, all" + echo + echo "Environment overrides:" + echo " CUVS_LUCENE_JAR, CUVS_LUCENE_CUVS_JAVA_JAR," + echo " CUVS_LUCENE_PYLUCENE_TEST_CLASSES," + echo " CUVS_LUCENE_PYLUCENE_CASES, CUVS_LUCENE_PYLUCENE_ROWS," + echo " CUVS_LUCENE_PYLUCENE_DIMS, CUVS_LUCENE_PYLUCENE_TOPK," + echo " CUVS_LUCENE_PYLUCENE_MIN_RECALL, PYTHON, MVN" + exit 0 + ;; + *) + echo "Unknown argument: ${arg}" >&2 + echo "Use -- before arguments intended for pytest." >&2 + exit 2 + ;; + esac +done + +require_command() { + if ! command -v "$1" >/dev/null 2>&1; then + echo "Required command not found: $1" >&2 + exit 127 + fi +} + +absolute_from_repo_root() { + case "$1" in + /*) + printf '%s\n' "$1" + ;; + *) + printf '%s/%s\n' "${REPO_ROOT}" "$1" + ;; + esac +} + +find_cuvs_java_jar() { + if [[ -n "${CUVS_LUCENE_CUVS_JAVA_JAR:-}" ]]; then + printf '%s\n' "${CUVS_LUCENE_CUVS_JAVA_JAR}" + return + fi + + local versioned_jar="${HOME}/.m2/repository/com/nvidia/cuvs/cuvs-java/${project_version}/cuvs-java-${project_version}.jar" + if [[ -f "${versioned_jar}" ]]; then + printf '%s\n' "${versioned_jar}" + return + fi + + local m2_repo="${HOME}/.m2/repository/com/nvidia/cuvs/cuvs-java" + if [[ ! -d "${m2_repo}" ]]; then + return + fi + + find "${m2_repo}" \ + -type f \ + -name 'cuvs-java-*.jar' \ + ! -name '*sources*' \ + ! -name '*javadoc*' \ + ! -name '*x86_64*' \ + | sort -V \ + | tail -n 1 +} + +require_command "${PYTHON_BIN}" + +if [[ "${SKIP_BUILD}" -eq 0 && -z "${CUVS_LUCENE_JAR_PATH}" ]]; then + require_command "${MVN_BIN}" + "${MVN_BIN}" clean package -DskipTests +fi + +project_version="$( + sed -n 's/.*CUVS_LUCENE#VERSION_UPDATE_MARKER_START-->\([^<]*\)<\/version>.*/\1/p' pom.xml \ + | head -n 1 +)" +if [[ -z "${project_version}" ]]; then + echo "Unable to determine project version from pom.xml" >&2 + exit 1 +fi + +if [[ -n "${CUVS_LUCENE_JAR_PATH}" ]]; then + cuvs_lucene_jar="${CUVS_LUCENE_JAR_PATH}" +else + cuvs_lucene_jar="target/cuvs-lucene-${project_version}.jar" +fi +cuvs_lucene_jar_abs="$(absolute_from_repo_root "${cuvs_lucene_jar}")" +if [[ ! -f "${cuvs_lucene_jar_abs}" ]]; then + echo "cuvs-lucene jar not found: ${cuvs_lucene_jar_abs}" >&2 + echo "Run without --no-build, or set CUVS_LUCENE_JAR." >&2 + exit 1 +fi + +cuvs_java_jar="$(find_cuvs_java_jar)" +if [[ -z "${cuvs_java_jar}" || ! -f "${cuvs_java_jar}" ]]; then + echo "Base cuvs-java jar not found." >&2 + echo "Set CUVS_LUCENE_CUVS_JAVA_JAR to the base jar, not a native classifier jar." >&2 + exit 1 +fi +cuvs_java_jar_abs="$(absolute_from_repo_root "${cuvs_java_jar}")" + +pylucene_test_classes_abs="$( + absolute_from_repo_root "${PYLUCENE_TEST_CLASSES_PATH}" +)" +if [[ ! -d "${pylucene_test_classes_abs}" ]]; then + echo "PyLucene test classes not found: ${pylucene_test_classes_abs}" >&2 + echo "Run Maven test compilation, or set CUVS_LUCENE_PYLUCENE_TEST_CLASSES." >&2 + exit 1 +fi +pylucene_test_classes_abs="$( + cd "${pylucene_test_classes_abs}" && pwd -P +)" + +"${PYTHON_BIN}" -c "import lucene" >/dev/null 2>&1 || { + echo "Python cannot import PyLucene's lucene module." >&2 + echo "Activate a PyLucene environment compatible with this project's Lucene version." >&2 + exit 1 +} + +"${PYTHON_BIN}" -m pytest --version >/dev/null 2>&1 || { + echo "Python cannot run pytest. Install pytest in the active PyLucene environment." >&2 + exit 1 +} + +pytest_env=( + "CUVS_LUCENE_JAR=${cuvs_lucene_jar_abs}" + "CUVS_LUCENE_CUVS_JAVA_JAR=${cuvs_java_jar_abs}" + "CUVS_LUCENE_PYLUCENE_TEST_CLASSES=${pylucene_test_classes_abs}" + "CUVS_LUCENE_VERIFY_ALL_CODECS=${CUVS_LUCENE_VERIFY_ALL_CODECS:-1}" +) + +if [[ "${FULL_END_TO_END}" -eq 1 ]]; then + pytest_env+=( + "CUVS_LUCENE_PYLUCENE_CASES=${PYLUCENE_CASES:-all}" + "CUVS_LUCENE_PYLUCENE_ROWS=${PYLUCENE_ROWS:-2000}" + "CUVS_LUCENE_PYLUCENE_DIMS=${PYLUCENE_DIMS:-32}" + "CUVS_LUCENE_PYLUCENE_TOPK=${PYLUCENE_TOPK:-20}" + ) +else + if [[ -n "${PYLUCENE_CASES}" ]]; then + pytest_env+=("CUVS_LUCENE_PYLUCENE_CASES=${PYLUCENE_CASES}") + fi + if [[ -n "${PYLUCENE_ROWS}" ]]; then + pytest_env+=("CUVS_LUCENE_PYLUCENE_ROWS=${PYLUCENE_ROWS}") + fi + if [[ -n "${PYLUCENE_DIMS}" ]]; then + pytest_env+=("CUVS_LUCENE_PYLUCENE_DIMS=${PYLUCENE_DIMS}") + fi + if [[ -n "${PYLUCENE_TOPK}" ]]; then + pytest_env+=("CUVS_LUCENE_PYLUCENE_TOPK=${PYLUCENE_TOPK}") + fi +fi + +if [[ -n "${PYLUCENE_MIN_RECALL}" ]]; then + pytest_env+=("CUVS_LUCENE_PYLUCENE_MIN_RECALL=${PYLUCENE_MIN_RECALL}") +fi + +env "${pytest_env[@]}" \ + "${PYTHON_BIN}" -m pytest -q -s \ + src/test/python/test_pylucene_end_to_end.py \ + "${PYTEST_ARGS[@]}" diff --git a/pom.xml b/pom.xml index 70f246d2..18afe68c 100644 --- a/pom.xml +++ b/pom.xml @@ -166,9 +166,9 @@ maven-assembly-plugin 3.6.0 - - jar-with-dependencies - + + src/main/assembly/jar-with-dependencies.xml + diff --git a/src/main/assembly/jar-with-dependencies.xml b/src/main/assembly/jar-with-dependencies.xml new file mode 100644 index 00000000..3722a871 --- /dev/null +++ b/src/main/assembly/jar-with-dependencies.xml @@ -0,0 +1,22 @@ + + jar-with-dependencies + + jar + + false + + + / + true + true + runtime + + + + + metaInf-services + + + diff --git a/src/main/java/com/nvidia/cuvs/lucene/CuVS2510GPUSearchCodec.java b/src/main/java/com/nvidia/cuvs/lucene/CuVS2510GPUSearchCodec.java index 2de761a7..2567fc4f 100644 --- a/src/main/java/com/nvidia/cuvs/lucene/CuVS2510GPUSearchCodec.java +++ b/src/main/java/com/nvidia/cuvs/lucene/CuVS2510GPUSearchCodec.java @@ -29,7 +29,7 @@ public class CuVS2510GPUSearchCodec extends FilterCodec { * @throws Exception */ public CuVS2510GPUSearchCodec() throws Exception { - this(NAME, LuceneProvider.getCodec("101")); + this(NAME, LuceneProvider.getDefaultDelegateCodec()); initializeFormat(new GPUSearchParams.Builder().build()); } @@ -53,7 +53,7 @@ public CuVS2510GPUSearchCodec(String name, Codec delegate) { * @throws Exception Exception raised when initializing the codec */ public CuVS2510GPUSearchCodec(GPUSearchParams params) throws Exception { - this(NAME, LuceneProvider.getCodec("101")); + this(NAME, LuceneProvider.getDefaultDelegateCodec()); initializeFormat(params); } diff --git a/src/main/java/com/nvidia/cuvs/lucene/CuVS2510GPUVectorsFormat.java b/src/main/java/com/nvidia/cuvs/lucene/CuVS2510GPUVectorsFormat.java index ccc61eae..a4caf665 100644 --- a/src/main/java/com/nvidia/cuvs/lucene/CuVS2510GPUVectorsFormat.java +++ b/src/main/java/com/nvidia/cuvs/lucene/CuVS2510GPUVectorsFormat.java @@ -77,6 +77,11 @@ public KnnVectorsWriter fieldsWriter(SegmentWriteState state) throws IOException return new CuVS2510GPUVectorsWriter(state, gpuSearchParams, flatWriter); } + @Override + public String toString() { + return getName() + "(" + WriterTelemetry.forCagra(gpuSearchParams) + ")"; + } + /** * Returns a KnnVectorsReader instance to read the vectors from the index. */ diff --git a/src/main/java/com/nvidia/cuvs/lucene/Lucene101AcceleratedHNSWBaseLayerCodec.java b/src/main/java/com/nvidia/cuvs/lucene/Lucene101AcceleratedHNSWBaseLayerCodec.java new file mode 100644 index 00000000..9fab1971 --- /dev/null +++ b/src/main/java/com/nvidia/cuvs/lucene/Lucene101AcceleratedHNSWBaseLayerCodec.java @@ -0,0 +1,31 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +import com.nvidia.cuvs.CagraIndexParams.CagraGraphBuildAlgo; + +/** + * Accelerated HNSW codec that builds an intermediate CAGRA graph and writes only the base HNSW layer. + */ +public class Lucene101AcceleratedHNSWBaseLayerCodec extends Lucene101AcceleratedHNSWCodec { + + private static final String NAME = "Lucene101AcceleratedHNSWBaseLayerCodec"; + private static final int CAGRA_GRAPH_DEGREE = 32; + private static final int CAGRA_INTERMEDIATE_GRAPH_DEGREE = 64; + + /** Default constructor used by Lucene SPI. */ + public Lucene101AcceleratedHNSWBaseLayerCodec() throws Exception { + super( + NAME, + LuceneProvider.getDefaultDelegateCodec(), + new AcceleratedHNSWParams.Builder() + .withStrategy(AcceleratedHNSWParams.Strategy.CUSTOM) + .withCagraGraphBuildAlgo(CagraGraphBuildAlgo.NN_DESCENT) + .withGraphDegree(CAGRA_GRAPH_DEGREE) + .withIntermediateGraphDegree(CAGRA_INTERMEDIATE_GRAPH_DEGREE) + .withHNSWLayer(1) + .build()); + } +} diff --git a/src/main/java/com/nvidia/cuvs/lucene/Lucene101AcceleratedHNSWCodec.java b/src/main/java/com/nvidia/cuvs/lucene/Lucene101AcceleratedHNSWCodec.java index b4c5a33d..7ab4a154 100644 --- a/src/main/java/com/nvidia/cuvs/lucene/Lucene101AcceleratedHNSWCodec.java +++ b/src/main/java/com/nvidia/cuvs/lucene/Lucene101AcceleratedHNSWCodec.java @@ -30,7 +30,7 @@ public class Lucene101AcceleratedHNSWCodec extends FilterCodec { * @throws Exception */ public Lucene101AcceleratedHNSWCodec() throws Exception { - this(NAME, LuceneProvider.getCodec("101")); + this(NAME, LuceneProvider.getDefaultDelegateCodec()); } /** @@ -52,7 +52,19 @@ public Lucene101AcceleratedHNSWCodec(String name, Codec delegate) { */ public Lucene101AcceleratedHNSWCodec(AcceleratedHNSWParams acceleratedHNSWParams) throws Exception { - this(NAME, LuceneProvider.getCodec("101")); + this(NAME, LuceneProvider.getDefaultDelegateCodec(), acceleratedHNSWParams); + } + + /** + * Constructor for subclasses that expose named accelerated HNSW configurations via SPI. + * + * @param name the codec's name + * @param delegate the delegate codec to filter + * @param acceleratedHNSWParams instance of {@link AcceleratedHNSWParams} + */ + protected Lucene101AcceleratedHNSWCodec( + String name, Codec delegate, AcceleratedHNSWParams acceleratedHNSWParams) { + super(name, delegate); initializeFormat(acceleratedHNSWParams); } diff --git a/src/main/java/com/nvidia/cuvs/lucene/Lucene101AcceleratedHNSWMultiLayerCodec.java b/src/main/java/com/nvidia/cuvs/lucene/Lucene101AcceleratedHNSWMultiLayerCodec.java new file mode 100644 index 00000000..3ad87b4b --- /dev/null +++ b/src/main/java/com/nvidia/cuvs/lucene/Lucene101AcceleratedHNSWMultiLayerCodec.java @@ -0,0 +1,31 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +import com.nvidia.cuvs.CagraIndexParams.CagraGraphBuildAlgo; + +/** + * Accelerated HNSW codec that builds an intermediate CAGRA graph and writes multiple HNSW layers. + */ +public class Lucene101AcceleratedHNSWMultiLayerCodec extends Lucene101AcceleratedHNSWCodec { + + private static final String NAME = "Lucene101AcceleratedHNSWMultiLayerCodec"; + private static final int CAGRA_GRAPH_DEGREE = 32; + private static final int CAGRA_INTERMEDIATE_GRAPH_DEGREE = 64; + + /** Default constructor used by Lucene SPI. */ + public Lucene101AcceleratedHNSWMultiLayerCodec() throws Exception { + super( + NAME, + LuceneProvider.getDefaultDelegateCodec(), + new AcceleratedHNSWParams.Builder() + .withStrategy(AcceleratedHNSWParams.Strategy.CUSTOM) + .withCagraGraphBuildAlgo(CagraGraphBuildAlgo.NN_DESCENT) + .withGraphDegree(CAGRA_GRAPH_DEGREE) + .withIntermediateGraphDegree(CAGRA_INTERMEDIATE_GRAPH_DEGREE) + .withHNSWLayer(3) + .build()); + } +} diff --git a/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsFormat.java b/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsFormat.java index 3c42707c..290a0586 100644 --- a/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsFormat.java +++ b/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsFormat.java @@ -4,6 +4,7 @@ */ package com.nvidia.cuvs.lucene; +import static com.nvidia.cuvs.lucene.ThreadLocalCuVSResourcesProvider.isCpuHnswFallbackForced; import static com.nvidia.cuvs.lucene.ThreadLocalCuVSResourcesProvider.isSupported; import com.nvidia.cuvs.LibraryException; @@ -79,9 +80,13 @@ public KnnVectorsWriter fieldsWriter(SegmentWriteState state) throws IOException log.log(Level.FINE, "cuVS is supported so using the Lucene99AcceleratedHNSWVectorsWriter"); return new Lucene99AcceleratedHNSWVectorsWriter(state, acceleratedHNSWParams, flatWriter); } else { + boolean forcedCpuFallback = isCpuHnswFallbackForced(); log.log( - Level.WARNING, - "GPU based indexing not supported, falling back to using the Lucene99HnswVectorsWriter"); + forcedCpuFallback ? Level.FINE : Level.WARNING, + forcedCpuFallback + ? "Forced CPU HNSW fallback, using the Lucene99HnswVectorsWriter" + : "GPU based indexing not supported, falling back to using the" + + " Lucene99HnswVectorsWriter"); try { return LUCENE_PROVIDER.getLuceneHnswVectorsWriterInstance( state, @@ -96,6 +101,11 @@ public KnnVectorsWriter fieldsWriter(SegmentWriteState state) throws IOException } } + @Override + public String toString() { + return getName() + "(" + WriterTelemetry.forHnsw(acceleratedHNSWParams) + ")"; + } + /** * Returns a KnnVectorsReader to read the vectors from the index. */ diff --git a/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedCodec.java b/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedCodec.java index f2c1aa37..5e149b49 100644 --- a/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedCodec.java +++ b/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedCodec.java @@ -25,7 +25,7 @@ public class LuceneAcceleratedHNSWBinaryQuantizedCodec extends FilterCodec { private KnnVectorsFormat format; public LuceneAcceleratedHNSWBinaryQuantizedCodec() throws Exception { - this(NAME, LuceneProvider.getCodec("101")); + this(NAME, LuceneProvider.getDefaultDelegateCodec()); } public LuceneAcceleratedHNSWBinaryQuantizedCodec(String name, Codec delegate) { @@ -35,7 +35,7 @@ public LuceneAcceleratedHNSWBinaryQuantizedCodec(String name, Codec delegate) { public LuceneAcceleratedHNSWBinaryQuantizedCodec(AcceleratedHNSWParams acceleratedHNSWParams) throws Exception { - this(NAME, LuceneProvider.getCodec("101")); + super(NAME, LuceneProvider.getDefaultDelegateCodec()); initializeFormat(acceleratedHNSWParams); } diff --git a/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat.java b/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat.java index 0f8d9602..c1ab9998 100644 --- a/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat.java +++ b/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat.java @@ -13,7 +13,6 @@ import org.apache.lucene.codecs.KnnVectorsFormat; import org.apache.lucene.codecs.KnnVectorsReader; import org.apache.lucene.codecs.KnnVectorsWriter; -import org.apache.lucene.codecs.hnsw.DefaultFlatVectorScorer; import org.apache.lucene.codecs.hnsw.FlatVectorsFormat; import org.apache.lucene.index.SegmentReadState; import org.apache.lucene.index.SegmentWriteState; @@ -27,24 +26,61 @@ public class LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat extends KnnVector private static final Logger log = Logger.getLogger(LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat.class.getName()); - private static final LuceneProvider LUCENE102_PROVIDER; - private static final LuceneProvider LUCENE99_PROVIDER; - private static final FlatVectorsFormat FLAT_VECTORS_FORMAT; private static final int MAX_DIMENSIONS = 4096; + private static final LuceneProvider LUCENE_99_PROVIDER = + getLuceneProvider(LuceneProvider.LUCENE_99_FORMAT_VERSION); + private static volatile FlatVectorsFormat cachedFlatVectorsFormat; private final AcceleratedHNSWParams acceleratedHNSWParams; + private volatile KnnVectorsFormat cachedFallbackFormat; - static { + private static LuceneProvider getLuceneProvider(String version) { try { - LUCENE99_PROVIDER = LuceneProvider.getInstance("99"); - LUCENE102_PROVIDER = LuceneProvider.getInstance("102"); - FLAT_VECTORS_FORMAT = - LUCENE102_PROVIDER.getLuceneFlatVectorsFormatInstance(DefaultFlatVectorScorer.INSTANCE); + return LuceneProvider.getInstance(version); } catch (Exception e) { - throw new ExceptionInInitializerError(e.getMessage()); + throw new UnsupportedOperationException( + "Lucene" + version + " vector formats are not available in this runtime", e); } } + private static FlatVectorsFormat getOrCreateFlatVectorsFormat() { + FlatVectorsFormat format = cachedFlatVectorsFormat; + if (format == null) { + synchronized (LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat.class) { + format = cachedFlatVectorsFormat; + if (format == null) { + try { + format = + getLuceneProvider(LuceneProvider.LUCENE_102_BINARY_FORMAT_VERSION) + .getLuceneBinaryQuantizedVectorsFormatInstance(); + cachedFlatVectorsFormat = format; + } catch (Exception e) { + throw new UnsupportedOperationException( + "Binary quantized vectors require Lucene102 vector formats", e); + } + } + } + } + return format; + } + + private KnnVectorsFormat getOrCreateFallbackFormat() throws Exception { + KnnVectorsFormat format = cachedFallbackFormat; + if (format == null) { + synchronized (this) { + format = cachedFallbackFormat; + if (format == null) { + format = + getLuceneProvider(LuceneProvider.LUCENE_102_BINARY_FORMAT_VERSION) + .getLuceneHnswBinaryQuantizedVectorsFormatInstance( + acceleratedHNSWParams.getMaxConn(), acceleratedHNSWParams.getBeamWidth()); + cachedFallbackFormat = format; + } + } + } + return format; + } + /** * Initializes {@link LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat} with default values. * @@ -70,8 +106,8 @@ public LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat( */ @Override public KnnVectorsWriter fieldsWriter(SegmentWriteState state) throws IOException { - var flatWriter = FLAT_VECTORS_FORMAT.fieldsWriter(state); if (isSupported()) { + var flatWriter = getOrCreateFlatVectorsFormat().fieldsWriter(state); log.log( Level.FINE, "cuVS is supported so using the Lucene99AcceleratedHNSWBinaryQuantizedVectorsWriter"); @@ -84,12 +120,9 @@ public KnnVectorsWriter fieldsWriter(SegmentWriteState state) throws IOException Level.WARNING, "GPU based indexing not supported, falling back to using the" + " Lucene102HnswBinaryQuantizedVectorsFormat"); - KnnVectorsFormat fallbackFormat = - LUCENE102_PROVIDER.getLuceneHnswBinaryQuantizedVectorsFormatInstance( - acceleratedHNSWParams.getMaxConn(), acceleratedHNSWParams.getBeamWidth()); - return fallbackFormat.fieldsWriter(state); + return getOrCreateFallbackFormat().fieldsWriter(state); } catch (Exception e) { - throw new RuntimeException(e.getMessage()); + throw new IOException("Unable to initialize the binary quantized fallback writer", e); } } } @@ -100,10 +133,10 @@ public KnnVectorsWriter fieldsWriter(SegmentWriteState state) throws IOException @Override public KnnVectorsReader fieldsReader(SegmentReadState state) throws IOException { try { - return LUCENE99_PROVIDER.getLuceneHnswVectorsReaderInstance( - state, FLAT_VECTORS_FORMAT.fieldsReader(state)); + return LUCENE_99_PROVIDER.getLuceneHnswVectorsReaderInstance( + state, getOrCreateFlatVectorsFormat().fieldsReader(state)); } catch (Exception e) { - throw new RuntimeException(e.getMessage()); + throw new IOException("Unable to initialize the binary quantized vectors reader", e); } } diff --git a/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedCodec.java b/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedCodec.java index 0c7736a0..4e49027c 100644 --- a/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedCodec.java +++ b/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedCodec.java @@ -25,7 +25,7 @@ public class LuceneAcceleratedHNSWScalarQuantizedCodec extends FilterCodec { private KnnVectorsFormat format; public LuceneAcceleratedHNSWScalarQuantizedCodec() throws Exception { - this(NAME, LuceneProvider.getCodec("101")); + this(NAME, LuceneProvider.getDefaultDelegateCodec()); } public LuceneAcceleratedHNSWScalarQuantizedCodec(String name, Codec delegate) { @@ -35,7 +35,7 @@ public LuceneAcceleratedHNSWScalarQuantizedCodec(String name, Codec delegate) { public LuceneAcceleratedHNSWScalarQuantizedCodec(AcceleratedHNSWParams acceleratedHNSWParams) throws Exception { - this(NAME, LuceneProvider.getCodec("101")); + this(NAME, LuceneProvider.getDefaultDelegateCodec()); initializeFormat(acceleratedHNSWParams); } diff --git a/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsFormat.java b/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsFormat.java index 534390b8..6c2ab076 100644 --- a/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsFormat.java +++ b/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsFormat.java @@ -65,8 +65,8 @@ public LuceneAcceleratedHNSWScalarQuantizedVectorsFormat( */ @Override public KnnVectorsWriter fieldsWriter(SegmentWriteState state) throws IOException { - var flatWriter = FLAT_VECTORS_FORMAT.fieldsWriter(state); if (isSupported()) { + var flatWriter = FLAT_VECTORS_FORMAT.fieldsWriter(state); log.info("cuVS is supported so using the Lucene99AcceleratedHNSWQuantizedVectorsWriter"); return new LuceneAcceleratedHNSWScalarQuantizedVectorsWriter( state, acceleratedHNSWParams, flatWriter); diff --git a/src/main/java/com/nvidia/cuvs/lucene/LuceneProvider.java b/src/main/java/com/nvidia/cuvs/lucene/LuceneProvider.java index 7635e323..7ecab96e 100644 --- a/src/main/java/com/nvidia/cuvs/lucene/LuceneProvider.java +++ b/src/main/java/com/nvidia/cuvs/lucene/LuceneProvider.java @@ -8,10 +8,14 @@ import java.lang.invoke.VarHandle; import java.lang.reflect.Constructor; import java.lang.reflect.InvocationTargetException; +import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.logging.Level; import java.util.logging.Logger; import org.apache.lucene.codecs.Codec; +import org.apache.lucene.codecs.KnnVectorsFormat; import org.apache.lucene.codecs.KnnVectorsReader; import org.apache.lucene.codecs.KnnVectorsWriter; import org.apache.lucene.codecs.hnsw.FlatVectorsFormat; @@ -31,6 +35,11 @@ public class LuceneProvider { static final Logger log = Logger.getLogger(LuceneProvider.class.getName()); + static final String LUCENE_99_FORMAT_VERSION = "99"; + static final String LUCENE_102_BINARY_FORMAT_VERSION = "102"; + + private static final List SUPPORTED_DELEGATE_CODEC_VERSIONS = + List.of("101", LUCENE_99_FORMAT_VERSION); private static final String BASE = "org.apache.lucene."; private static String codecs = "codecs.lucene."; @@ -79,7 +88,7 @@ public class LuceneProvider { private static String luceneCodec = BASE + codecs + "LuceneCodec"; private static String luceneCodecFallback = BASE + fallbackCodecs + "LuceneCodec"; - private static LuceneProvider instance; + private static final Map INSTANCES = new HashMap<>(); private static MethodHandles.Lookup lookup = MethodHandles.lookup(); @@ -92,14 +101,30 @@ public class LuceneProvider { private Class scalarQuantizedVectorsFormat; private Class hnswScalarQuantizedVectorsFormat; - public static LuceneProvider getInstance(String version) throws ClassNotFoundException { + public static synchronized LuceneProvider getInstance(String version) + throws ClassNotFoundException { + LuceneProvider instance = INSTANCES.get(version); if (instance == null) { instance = new LuceneProvider(version); + INSTANCES.put(version, instance); } return instance; } private LuceneProvider(String version) throws ClassNotFoundException { + // TODO: Find a better way if possible, but as a separate initiative. + if (LUCENE_102_BINARY_FORMAT_VERSION.equals(version)) { + binaryQuantizedVectorsFormat = + loadClass( + setVersion(luceneBinaryQuantizedVectorsFormat, version), + setVersion(luceneBinaryQuantizedVectorsFormatFallback, version)); + hnswBinaryQuantizedVectorsFormat = + loadClass( + setVersion(luceneHnswBinaryQuantizedVectorsFormat, version), + setVersion(luceneHnswBinaryQuantizedVectorsFormatFallback, version)); + return; + } + flatVectorsFormat = loadClass( setVersion(luceneFlatVectorsFormat, version), @@ -125,18 +150,6 @@ private LuceneProvider(String version) throws ClassNotFoundException { loadClass( setVersion(luceneHnswScalarQuantizedVectorsFormat, version), setVersion(luceneHnswScalarQuantizedVectorsFormatFallback, version)); - - // TODO: Find a better way if possible, but as a separate initiative. - if ("102".equals(version)) { - binaryQuantizedVectorsFormat = - loadClass( - setVersion(luceneBinaryQuantizedVectorsFormat, version), - setVersion(luceneBinaryQuantizedVectorsFormatFallback, version)); - hnswBinaryQuantizedVectorsFormat = - loadClass( - setVersion(luceneHnswBinaryQuantizedVectorsFormat, version), - setVersion(luceneHnswBinaryQuantizedVectorsFormatFallback, version)); - } } private static String setVersion(String pkg, String version) { @@ -147,14 +160,20 @@ private static Class loadClass(String defaultClassName, String fallbackClassN throws ClassNotFoundException { try { return Class.forName(defaultClassName); - } catch (ClassNotFoundException e) { + } catch (ClassNotFoundException defaultException) { // Load class from fallback package. try { return Class.forName(fallbackClassName); - } catch (ClassNotFoundException e1) { - // Should not reach here. - log.log(Level.SEVERE, "Unable to load class: " + fallbackClassName); - throw e1; + } catch (ClassNotFoundException fallbackException) { + ClassNotFoundException missing = + new ClassNotFoundException( + "Unable to load Lucene class. Tried " + + defaultClassName + + " and " + + fallbackClassName); + missing.addSuppressed(defaultException); + missing.addSuppressed(fallbackException); + throw missing; } } } @@ -173,6 +192,26 @@ public static Codec getCodec(String version) return (Codec) codecClassConstructor.newInstance(); } + public static Codec getDefaultDelegateCodec() { + List failures = new ArrayList<>(); + for (String version : SUPPORTED_DELEGATE_CODEC_VERSIONS) { + try { + return getCodec(version); + } catch (ReflectiveOperationException + | SecurityException + | IllegalArgumentException + | LinkageError e) { + failures.add("Lucene" + version + ": " + e.getMessage()); + log.log(Level.FINE, "Unable to load Lucene" + version + "Codec", e); + } + } + throw new IllegalStateException( + "Unable to load a supported Lucene delegate codec. Tried " + + SUPPORTED_DELEGATE_CODEC_VERSIONS + + ". Failures: " + + failures); + } + public FlatVectorsFormat getLuceneFlatVectorsFormatInstance(FlatVectorsScorer scorer) throws Exception { try { @@ -245,7 +284,7 @@ public List getSimilarityFunctions() } } - public FlatVectorsFormat getluceneBinaryQuantizedVectorsFormatInstance() throws Exception { + public FlatVectorsFormat getLuceneBinaryQuantizedVectorsFormatInstance() throws Exception { try { Constructor luceneBinaryQuantizedVectorsFormatConstructor = binaryQuantizedVectorsFormat.getConstructor(); @@ -258,17 +297,17 @@ public FlatVectorsFormat getluceneBinaryQuantizedVectorsFormatInstance() throws } } - public FlatVectorsFormat getLuceneHnswBinaryQuantizedVectorsFormatInstance( + public KnnVectorsFormat getLuceneHnswBinaryQuantizedVectorsFormatInstance( int maxConn, int beamWidth) throws Exception { try { Constructor luceneHnswBinaryQuantizedVectorsFormatConstructor = - hnswBinaryQuantizedVectorsFormat.getConstructor(Integer.TYPE, Integer.TYPE); - return (FlatVectorsFormat) + hnswBinaryQuantizedVectorsFormat.getConstructor(int.class, int.class); + return (KnnVectorsFormat) luceneHnswBinaryQuantizedVectorsFormatConstructor.newInstance(maxConn, beamWidth); } catch (Exception e) { log.log( Level.SEVERE, - "Unable to initialize LuceneBinaryQuantizedVectorsFormat: " + e.getMessage()); + "Unable to initialize LuceneHnswBinaryQuantizedVectorsFormat: " + e.getMessage()); throw e; } } diff --git a/src/main/java/com/nvidia/cuvs/lucene/ThreadLocalCuVSResourcesProvider.java b/src/main/java/com/nvidia/cuvs/lucene/ThreadLocalCuVSResourcesProvider.java index 9e259e27..a13271d4 100644 --- a/src/main/java/com/nvidia/cuvs/lucene/ThreadLocalCuVSResourcesProvider.java +++ b/src/main/java/com/nvidia/cuvs/lucene/ThreadLocalCuVSResourcesProvider.java @@ -18,6 +18,7 @@ public class ThreadLocalCuVSResourcesProvider { private static final Logger log = Logger.getLogger(ThreadLocalCuVSResourcesProvider.class.getName()); + static final String FORCE_CPU_HNSW_FALLBACK_PROPERTY = "cuvs.lucene.forceCpuHnswFallback"; private static final ThreadLocal cuVSResources; static { @@ -30,6 +31,9 @@ public class ThreadLocalCuVSResourcesProvider { * @return an instance of CuVSResources */ public static CuVSResources getCuVSResourcesInstance() { + if (isCpuHnswFallbackForced()) { + return null; + } return cuVSResources.get(); } @@ -75,7 +79,7 @@ public static void closeCuVSResourcesInstance() { * @throws UnsupportedOperationException */ public static void assertIsSupported() throws UnsupportedOperationException { - if (cuVSResources.get() == null) { + if (isCpuHnswFallbackForced() || cuVSResources.get() == null) { throw new UnsupportedOperationException("cuVS is not supported"); } } @@ -86,6 +90,10 @@ public static void assertIsSupported() throws UnsupportedOperationException { * @return true if cuVS is supported else false */ public static boolean isSupported() { - return cuVSResources.get() != null; + return !isCpuHnswFallbackForced() && cuVSResources.get() != null; + } + + static boolean isCpuHnswFallbackForced() { + return Boolean.getBoolean(FORCE_CPU_HNSW_FALLBACK_PROPERTY); } } diff --git a/src/main/java/com/nvidia/cuvs/lucene/WriterTelemetry.java b/src/main/java/com/nvidia/cuvs/lucene/WriterTelemetry.java new file mode 100644 index 00000000..5095e8c2 --- /dev/null +++ b/src/main/java/com/nvidia/cuvs/lucene/WriterTelemetry.java @@ -0,0 +1,45 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +/** Formats vector-writer configuration for diagnostic output. */ +final class WriterTelemetry { + + private WriterTelemetry() {} + + static String forCagra(GPUSearchParams params) { + return "writerPath=gpu-cagra" + + ";cagraStrategy=" + + params.getStrategy().name() + + ";cagraGraphBuildAlgo=" + + params.getCagraGraphBuildAlgo().name() + + ";cagraGraphDegree=" + + params.getGraphdegree() + + ";cagraIntermediateGraphDegree=" + + params.getIntermediateGraphDegree(); + } + + static String forHnsw(AcceleratedHNSWParams params) { + String writerPath; + if (ThreadLocalCuVSResourcesProvider.isSupported()) { + writerPath = "gpu-hnsw"; + } else if (ThreadLocalCuVSResourcesProvider.isCpuHnswFallbackForced()) { + writerPath = "cpu-hnsw-fallback"; + } else { + writerPath = "cpu-hnsw-auto-fallback"; + } + + return "writerPath=" + + writerPath + + ";hnswLayers=" + + params.getHnswLayers() + + ";cagraGraphBuildAlgo=" + + params.getCagraGraphBuildAlgo().name() + + ";cagraGraphDegree=" + + params.getGraphdegree() + + ";cagraIntermediateGraphDegree=" + + params.getIntermediateGraphDegree(); + } +} diff --git a/src/main/resources/META-INF/services/org.apache.lucene.codecs.Codec b/src/main/resources/META-INF/services/org.apache.lucene.codecs.Codec index faa0684c..2b9d6a7f 100644 --- a/src/main/resources/META-INF/services/org.apache.lucene.codecs.Codec +++ b/src/main/resources/META-INF/services/org.apache.lucene.codecs.Codec @@ -2,6 +2,8 @@ # SPDX-License-Identifier: Apache-2.0 com.nvidia.cuvs.lucene.Lucene101AcceleratedHNSWCodec +com.nvidia.cuvs.lucene.Lucene101AcceleratedHNSWBaseLayerCodec +com.nvidia.cuvs.lucene.Lucene101AcceleratedHNSWMultiLayerCodec com.nvidia.cuvs.lucene.CuVS2510GPUSearchCodec com.nvidia.cuvs.lucene.LuceneAcceleratedHNSWBinaryQuantizedCodec com.nvidia.cuvs.lucene.LuceneAcceleratedHNSWScalarQuantizedCodec diff --git a/src/main/resources/META-INF/services/org.apache.lucene.codecs.KnnVectorsFormat b/src/main/resources/META-INF/services/org.apache.lucene.codecs.KnnVectorsFormat index 6625ac72..1f9ceeda 100644 --- a/src/main/resources/META-INF/services/org.apache.lucene.codecs.KnnVectorsFormat +++ b/src/main/resources/META-INF/services/org.apache.lucene.codecs.KnnVectorsFormat @@ -1,8 +1,6 @@ # SPDX-FileCopyrightText: Copyright (c) 2025-2026, NVIDIA CORPORATION. # SPDX-License-Identifier: Apache-2.0 -org.apache.lucene.codecs.lucene99.Lucene99HnswVectorsFormat -org.apache.lucene.codecs.lucene99.Lucene99HnswScalarQuantizedVectorsFormat com.nvidia.cuvs.lucene.CuVS2510GPUVectorsFormat com.nvidia.cuvs.lucene.Lucene99AcceleratedHNSWVectorsFormat com.nvidia.cuvs.lucene.LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat diff --git a/src/test/java/com/nvidia/cuvs/lucene/PyLuceneTestSupport.java b/src/test/java/com/nvidia/cuvs/lucene/PyLuceneTestSupport.java new file mode 100644 index 00000000..9350c714 --- /dev/null +++ b/src/test/java/com/nvidia/cuvs/lucene/PyLuceneTestSupport.java @@ -0,0 +1,382 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +import com.nvidia.cuvs.CagraIndexParams.CagraGraphBuildAlgo; +import java.io.IOException; +import java.lang.reflect.Field; +import java.lang.reflect.Method; +import java.util.Map; +import org.apache.lucene.codecs.KnnVectorsReader; +import org.apache.lucene.codecs.hnsw.HnswGraphProvider; +import org.apache.lucene.codecs.perfield.PerFieldKnnVectorsFormat; +import org.apache.lucene.index.CodecReader; +import org.apache.lucene.index.FilterLeafReader; +import org.apache.lucene.index.LeafReader; +import org.apache.lucene.index.LeafReaderContext; +import org.apache.lucene.index.QueryTimeout; +import org.apache.lucene.index.Term; +import org.apache.lucene.search.DocIdSetIterator; +import org.apache.lucene.search.KnnFloatVectorQuery; +import org.apache.lucene.search.Query; +import org.apache.lucene.search.TermQuery; +import org.apache.lucene.search.TopDocs; +import org.apache.lucene.search.knn.KnnCollectorManager; +import org.apache.lucene.util.Bits; + +/** + * Test-only adapters that PyLucene creates reflectively. Public no-argument constructors read query + * configuration from Java system properties. + */ +public final class PyLuceneTestSupport { + + public static final String QUERY_FIELD_PROPERTY = "cuvs.lucene.pylucene.query.field"; + public static final String QUERY_TARGET_PROPERTY = "cuvs.lucene.pylucene.query.target"; + public static final String QUERY_K_PROPERTY = "cuvs.lucene.pylucene.query.k"; + public static final String QUERY_I_TOP_K_PROPERTY = "cuvs.lucene.pylucene.query.iTopK"; + public static final String QUERY_SEARCH_WIDTH_PROPERTY = "cuvs.lucene.pylucene.query.searchWidth"; + public static final String QUERY_EXPECTED_HNSW_M_PROPERTY = + "cuvs.lucene.pylucene.query.expectedHnswM"; + public static final String QUERY_FILTER_FIELD_PROPERTY = "cuvs.lucene.pylucene.query.filterField"; + public static final String QUERY_FILTER_VALUE_PROPERTY = "cuvs.lucene.pylucene.query.filterValue"; + + private static final int CAGRA_GRAPH_DEGREE = 32; + private static final int CAGRA_INTERMEDIATE_GRAPH_DEGREE = 64; + + private PyLuceneTestSupport() {} + + /** CAGRA codec with graph degree 32 and intermediate graph degree 64. */ + public static final class CagraSearchCodec extends CuVS2510GPUSearchCodec { + + public CagraSearchCodec() throws Exception { + super( + new GPUSearchParams.Builder() + .withStrategy(GPUSearchParams.Strategy.CUSTOM) + .withCagraGraphBuildAlgo(CagraGraphBuildAlgo.NN_DESCENT) + .withGraphDegree(CAGRA_GRAPH_DEGREE) + .withIntermediateGraphDegree(CAGRA_INTERMEDIATE_GRAPH_DEGREE) + .build()); + } + } + + public static final class HnswGraphVerifyingQuery extends KnnFloatVectorQuery { + + private final int expectedM; + + public HnswGraphVerifyingQuery() { + this(QueryProperties.fromSystemPropertiesForHnsw()); + } + + private HnswGraphVerifyingQuery(QueryProperties properties) { + super(properties.field, properties.target, properties.k, properties.filter); + expectedM = parsePositiveInt(QUERY_EXPECTED_HNSW_M_PROPERTY); + } + + @Override + protected TopDocs approximateSearch( + LeafReaderContext context, + Bits acceptDocs, + int visitedLimit, + KnnCollectorManager knnCollectorManager) + throws IOException { + requireExpectedHnswGraph(context); + return super.approximateSearch(context, acceptDocs, visitedLimit, knnCollectorManager); + } + + @Override + protected TopDocs exactSearch( + LeafReaderContext context, DocIdSetIterator acceptIterator, QueryTimeout queryTimeout) + throws IOException { + requireExpectedHnswGraph(context); + return super.exactSearch(context, acceptIterator, queryTimeout); + } + + private void requireExpectedHnswGraph(LeafReaderContext context) { + KnnVectorsReader vectorsReader = vectorReaderForField(context, getField()); + if (!(vectorsReader instanceof HnswGraphProvider)) { + throw new AssertionError( + "HNSW graph-verifying query requires HnswGraphProvider for field '" + + getField() + + "', got " + + className(vectorsReader)); + } + + int actualM = persistedHnswM(vectorsReader, getField()); + if (actualM != expectedM) { + throw new AssertionError( + "HNSW graph-verifying query expected persisted M " + + expectedM + + " for field '" + + getField() + + "', got " + + actualM); + } + } + + private static int persistedHnswM(KnnVectorsReader vectorsReader, String fieldName) { + // The PyLucene bindings do not expose persisted M through their wrapped HNSW graph API. + try { + Field fieldsField = vectorsReader.getClass().getDeclaredField("fields"); + fieldsField.setAccessible(true); + Object fieldsValue = fieldsField.get(vectorsReader); + if (!(fieldsValue instanceof Map fields)) { + throw new AssertionError( + "Unexpected HNSW field metadata container: " + className(fieldsValue)); + } + Object fieldEntry = fields.get(fieldName); + if (fieldEntry == null) { + throw new AssertionError("Persisted HNSW metadata has no entry for field: " + fieldName); + } + Method mAccessor = fieldEntry.getClass().getDeclaredMethod("M"); + mAccessor.setAccessible(true); + Object value = mAccessor.invoke(fieldEntry); + if (!(value instanceof Integer persistedM)) { + throw new AssertionError("Unexpected persisted HNSW M value: " + value); + } + return persistedM; + } catch (ReflectiveOperationException | RuntimeException exception) { + throw new AssertionError( + "Unable to inspect persisted HNSW M for field '" + fieldName + "'", exception); + } + } + } + + public static final class CagraSearchQuery extends GPUKnnFloatVectorQuery { + + public CagraSearchQuery() { + this(QueryProperties.fromSystemProperties()); + } + + private CagraSearchQuery(QueryProperties properties) { + super( + properties.field, + properties.target, + properties.k, + properties.filter, + properties.iTopK, + properties.searchWidth); + } + + @Override + protected TopDocs approximateSearch( + LeafReaderContext context, + Bits acceptDocs, + int visitedLimit, + KnnCollectorManager knnCollectorManager) + throws IOException { + requireCagraOnlyIndex(context); + return super.approximateSearch(context, acceptDocs, visitedLimit, knnCollectorManager); + } + + @Override + protected TopDocs exactSearch( + LeafReaderContext context, DocIdSetIterator acceptIterator, QueryTimeout queryTimeout) { + throw new AssertionError("CAGRA search query must not use Lucene exact vector scoring"); + } + + private void requireCagraOnlyIndex(LeafReaderContext context) { + KnnVectorsReader vectorsReader = vectorReaderForField(context, getField()); + if (!(vectorsReader instanceof CuVS2510GPUVectorsReader cuvsReader)) { + throw new AssertionError( + "CAGRA search query requires CuVS2510GPUVectorsReader for field '" + + getField() + + "', got " + + className(vectorsReader)); + } + + var fieldInfo = cuvsReader.getFieldInfos().fieldInfo(getField()); + if (fieldInfo == null) { + throw new AssertionError( + "CAGRA search field is absent from the cuVS reader: " + getField()); + } + var cuvsIndexes = cuvsReader.getCuvsIndexes(); + if (cuvsIndexes == null) { + throw new AssertionError( + "CAGRA search found no loaded GPU indexes for field: " + getField()); + } + GPUIndex gpuIndex = cuvsIndexes.get(fieldInfo.number); + if (gpuIndex == null) { + throw new AssertionError("CAGRA search found no GPU index for field: " + getField()); + } + var cagraIndex = gpuIndex.getCagraIndex(); + if (cagraIndex == null) { + throw new AssertionError( + "CAGRA search requires a CAGRA index for field '" + + getField() + + "'; a brute-force or CPU fallback is not allowed"); + } + if (gpuIndex.getBruteforceIndex() != null) { + throw new AssertionError( + "CAGRA search requires a CAGRA-only index for field '" + + getField() + + "'; a brute-force index was also loaded"); + } + + long actualGraphDegree = cagraIndex.getGraph().columns(); + if (actualGraphDegree != CAGRA_GRAPH_DEGREE) { + throw new AssertionError( + "CAGRA search expected actual graph degree " + + CAGRA_GRAPH_DEGREE + + " for field '" + + getField() + + "', got " + + actualGraphDegree); + } + } + } + + private static KnnVectorsReader vectorReaderForField( + LeafReaderContext context, String fieldName) { + LeafReader leafReader = FilterLeafReader.unwrap(context.reader()); + if (!(leafReader instanceof CodecReader codecReader)) { + throw new AssertionError( + "Graph-verifying query requires a CodecReader leaf, got " + + leafReader.getClass().getName()); + } + + KnnVectorsReader vectorsReader = codecReader.getVectorReader(); + if (vectorsReader instanceof PerFieldKnnVectorsFormat.FieldsReader fieldsReader) { + vectorsReader = fieldsReader.getFieldReader(fieldName); + } + return vectorsReader; + } + + private static String className(Object value) { + return value == null ? "" : value.getClass().getName(); + } + + private static final class QueryProperties { + + private final String field; + private final float[] target; + private final int k; + private final int iTopK; + private final int searchWidth; + private final Query filter; + + private QueryProperties( + String field, float[] target, int k, int iTopK, int searchWidth, Query filter) { + this.field = field; + this.target = target; + this.k = k; + this.iTopK = iTopK; + this.searchWidth = searchWidth; + this.filter = filter; + } + + private static QueryProperties fromSystemProperties() { + String field = requiredProperty(QUERY_FIELD_PROPERTY); + float[] target = parseTarget(requiredProperty(QUERY_TARGET_PROPERTY)); + int k = parsePositiveInt(QUERY_K_PROPERTY); + int iTopK = parsePositiveInt(QUERY_I_TOP_K_PROPERTY); + int searchWidth = parsePositiveInt(QUERY_SEARCH_WIDTH_PROPERTY); + Query filter = optionalFilterFromSystemProperties(); + if (iTopK < k) { + throw new IllegalArgumentException( + "Java system property " + + QUERY_I_TOP_K_PROPERTY + + " must be greater than or equal to " + + QUERY_K_PROPERTY + + " (" + + k + + "), got: " + + iTopK); + } + return new QueryProperties(field, target, k, iTopK, searchWidth, filter); + } + + private static QueryProperties fromSystemPropertiesForHnsw() { + String field = requiredProperty(QUERY_FIELD_PROPERTY); + float[] target = parseTarget(requiredProperty(QUERY_TARGET_PROPERTY)); + int k = parsePositiveInt(QUERY_K_PROPERTY); + Query filter = optionalFilterFromSystemProperties(); + return new QueryProperties(field, target, k, k, 1, filter); + } + } + + private static Query optionalFilterFromSystemProperties() { + String field = System.getProperty(QUERY_FILTER_FIELD_PROPERTY); + String value = System.getProperty(QUERY_FILTER_VALUE_PROPERTY); + if (field == null && value == null) { + return null; + } + if (field == null || value == null) { + throw new IllegalArgumentException( + "Java system properties " + + QUERY_FILTER_FIELD_PROPERTY + + " and " + + QUERY_FILTER_VALUE_PROPERTY + + " must be set together"); + } + return new TermQuery( + new Term( + requiredProperty(QUERY_FILTER_FIELD_PROPERTY), + requiredProperty(QUERY_FILTER_VALUE_PROPERTY))); + } + + private static String requiredProperty(String propertyName) { + String value = System.getProperty(propertyName); + if (value == null || value.trim().isEmpty()) { + throw new IllegalArgumentException("Missing required Java system property: " + propertyName); + } + return value.trim(); + } + + private static int parsePositiveInt(String propertyName) { + String value = requiredProperty(propertyName); + final int parsed; + try { + parsed = Integer.parseInt(value); + } catch (NumberFormatException exception) { + throw new IllegalArgumentException( + "Java system property " + + propertyName + + " must be a positive integer, got: '" + + value + + "'", + exception); + } + if (parsed <= 0) { + throw new IllegalArgumentException( + "Java system property " + propertyName + " must be a positive integer, got: " + parsed); + } + return parsed; + } + + private static float[] parseTarget(String value) { + String[] values = value.split(",", -1); + float[] target = new float[values.length]; + for (int index = 0; index < values.length; index++) { + String component = values[index].trim(); + if (component.isEmpty()) { + throw malformedTarget(value, index, null); + } + try { + target[index] = Float.parseFloat(component); + } catch (NumberFormatException exception) { + throw malformedTarget(value, index, exception); + } + if (!Float.isFinite(target[index])) { + throw malformedTarget(value, index, null); + } + } + return target; + } + + private static IllegalArgumentException malformedTarget( + String value, int componentIndex, NumberFormatException cause) { + String message = + "Java system property " + + QUERY_TARGET_PROPERTY + + " must be a comma-separated list of finite floats; invalid component " + + componentIndex + + " in: '" + + value + + "'"; + return cause == null + ? new IllegalArgumentException(message) + : new IllegalArgumentException(message, cause); + } +} diff --git a/src/test/java/com/nvidia/cuvs/lucene/TestBackCompat.java b/src/test/java/com/nvidia/cuvs/lucene/TestBackCompat.java index 2de6e660..4b152ab0 100644 --- a/src/test/java/com/nvidia/cuvs/lucene/TestBackCompat.java +++ b/src/test/java/com/nvidia/cuvs/lucene/TestBackCompat.java @@ -6,8 +6,11 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; +import java.util.Set; import org.apache.lucene.codecs.Codec; import org.apache.lucene.codecs.hnsw.FlatVectorsFormat; import org.junit.Test; @@ -33,9 +36,50 @@ public void testNonexistentCodec() throws Exception { @Test public void testExistingComponents() throws Exception { - LuceneProvider provider = LuceneProvider.getInstance("99"); + LuceneProvider provider = LuceneProvider.getInstance(LuceneProvider.LUCENE_99_FORMAT_VERSION); assertTrue(provider.getLuceneFlatVectorsFormatInstance(null) instanceof FlatVectorsFormat); assertEquals(provider.getStaticIntParam("VERSION_CURRENT"), 0); assertNotEquals(provider.getSimilarityFunctions().size(), 0); } + + @Test + public void testProviderCachesSupportedVersion() throws Exception { + LuceneProvider lucene99Provider = + LuceneProvider.getInstance(LuceneProvider.LUCENE_99_FORMAT_VERSION); + assertSame( + lucene99Provider, LuceneProvider.getInstance(LuceneProvider.LUCENE_99_FORMAT_VERSION)); + } + + @Test + public void testProviderSupportsLucene102BinaryFormats() throws Exception { + LuceneProvider lucene102BinaryFormatProvider = + LuceneProvider.getInstance(LuceneProvider.LUCENE_102_BINARY_FORMAT_VERSION); + assertNotNull(lucene102BinaryFormatProvider.getLuceneBinaryQuantizedVectorsFormatInstance()); + assertNotNull( + lucene102BinaryFormatProvider.getLuceneHnswBinaryQuantizedVectorsFormatInstance(16, 100)); + } + + @Test + public void testDefaultDelegateCodec() { + Codec delegate = LuceneProvider.getDefaultDelegateCodec(); + assertNotNull(delegate); + assertTrue(Set.of("Lucene101", "Lucene99").contains(delegate.getName())); + assertTrue(delegate.getClass().getName().startsWith("org.apache.lucene.")); + } + + @Test + public void testServiceLoadedCodecsCanBeInstantiated() { + String[] codecNames = { + "Lucene101AcceleratedHNSWCodec", + "Lucene101AcceleratedHNSWBaseLayerCodec", + "Lucene101AcceleratedHNSWMultiLayerCodec", + "CuVS2510GPUSearchCodec", + "Lucene101AcceleratedHNSWBinaryQuantizedCodec", + "Lucene101AcceleratedHNSWScalarQuantizedCodec" + }; + for (String codecName : codecNames) { + assertTrue(Codec.availableCodecs().contains(codecName)); + assertEquals(codecName, Codec.forName(codecName).getName()); + } + } } diff --git a/src/test/java/com/nvidia/cuvs/lucene/TestWriterTelemetry.java b/src/test/java/com/nvidia/cuvs/lucene/TestWriterTelemetry.java new file mode 100644 index 00000000..b9f627ea --- /dev/null +++ b/src/test/java/com/nvidia/cuvs/lucene/TestWriterTelemetry.java @@ -0,0 +1,73 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNull; + +import com.nvidia.cuvs.CagraIndexParams.CagraGraphBuildAlgo; +import org.junit.Test; + +public class TestWriterTelemetry { + + private static final String FORCE_CPU_HNSW_FALLBACK_PROPERTY = "cuvs.lucene.forceCpuHnswFallback"; + + @Test + public void testCagraTelemetryIsComputedOnDemand() { + GPUSearchParams params = + new GPUSearchParams.Builder() + .withStrategy(GPUSearchParams.Strategy.CUSTOM) + .withCagraGraphBuildAlgo(CagraGraphBuildAlgo.NN_DESCENT) + .withGraphDegree(32) + .withIntermediateGraphDegree(64) + .build(); + + assertEquals( + "writerPath=gpu-cagra;cagraStrategy=CUSTOM;" + + "cagraGraphBuildAlgo=NN_DESCENT;cagraGraphDegree=32;" + + "cagraIntermediateGraphDegree=64", + WriterTelemetry.forCagra(params)); + assertEquals( + "CuVS2510GPUVectorsFormat(writerPath=gpu-cagra;cagraStrategy=CUSTOM;" + + "cagraGraphBuildAlgo=NN_DESCENT;cagraGraphDegree=32;" + + "cagraIntermediateGraphDegree=64)", + new CuVS2510GPUVectorsFormat(params).toString()); + assertNull(System.getProperty("cuvs.lucene.lastCagraWriterPath")); + } + + @Test + public void testForcedCpuHnswTelemetryIsComputedOnDemand() { + String previousValue = System.getProperty(FORCE_CPU_HNSW_FALLBACK_PROPERTY); + System.setProperty(FORCE_CPU_HNSW_FALLBACK_PROPERTY, "true"); + try { + AcceleratedHNSWParams params = + new AcceleratedHNSWParams.Builder() + .withHNSWLayer(3) + .withGraphDegree(32) + .withIntermediateGraphDegree(64) + .build(); + + assertEquals( + "writerPath=cpu-hnsw-fallback;hnswLayers=3;" + + "cagraGraphBuildAlgo=NN_DESCENT;cagraGraphDegree=32;" + + "cagraIntermediateGraphDegree=64", + WriterTelemetry.forHnsw(params)); + assertEquals( + "Lucene99AcceleratedHNSWVectorsFormat(" + + "writerPath=cpu-hnsw-fallback;hnswLayers=3;" + + "cagraGraphBuildAlgo=NN_DESCENT;cagraGraphDegree=32;" + + "cagraIntermediateGraphDegree=64)", + new Lucene99AcceleratedHNSWVectorsFormat(params).toString()); + assertNull(System.getProperty("cuvs.lucene.lastHnswWriterPath")); + assertNull(System.getProperty("cuvs.lucene.lastHnswLayers")); + } finally { + if (previousValue == null) { + System.clearProperty(FORCE_CPU_HNSW_FALLBACK_PROPERTY); + } else { + System.setProperty(FORCE_CPU_HNSW_FALLBACK_PROPERTY, previousValue); + } + } + } +} diff --git a/src/test/python/conftest.py b/src/test/python/conftest.py new file mode 100644 index 00000000..a56bfcf4 --- /dev/null +++ b/src/test/python/conftest.py @@ -0,0 +1,73 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +from __future__ import annotations + +import os + +import pytest + + +CASE_MARKER = "pylucene_case" +CASE_ENVIRONMENT_VARIABLE = "CUVS_LUCENE_PYLUCENE_CASES" +DEFAULT_CASE_SELECTION = "cpu-hnsw-1-segment" + + +def pytest_configure(config: pytest.Config) -> None: + config.addinivalue_line( + "markers", + f"{CASE_MARKER}(*selectors): select a PyLucene end-to-end scenario", + ) + + +def _requested_selectors() -> tuple[frozenset[str], bool]: + configured_value = os.environ.get(CASE_ENVIRONMENT_VARIABLE) + configured = configured_value or DEFAULT_CASE_SELECTION + selectors = frozenset( + selector.strip() + for selector in configured.split(",") + if selector.strip() + ) + if not selectors: + raise pytest.UsageError( + f"{CASE_ENVIRONMENT_VARIABLE} did not select any cases" + ) + return selectors, configured_value is not None + + +def pytest_collection_modifyitems( + config: pytest.Config, items: list[pytest.Item] +) -> None: + requested, explicitly_configured = _requested_selectors() + selectable_items: list[tuple[pytest.Item, frozenset[str]]] = [] + available: set[str] = set() + + for item in items: + marker = item.get_closest_marker(CASE_MARKER) + if marker is None: + continue + selectors = frozenset(str(value) for value in marker.args) + selectable_items.append((item, selectors)) + available.update(selectors) + + if not selectable_items: + return + + unknown = requested - available + if unknown and not explicitly_configured: + return + if unknown: + raise pytest.UsageError( + "Unknown PyLucene case or group: " + + ", ".join(sorted(unknown)) + ) + + deselected = [ + item + for item, selectors in selectable_items + if requested.isdisjoint(selectors) + ] + if deselected: + config.hook.pytest_deselected(items=deselected) + deselected_set = set(deselected) + items[:] = [item for item in items if item not in deselected_set] diff --git a/src/test/python/pylucene_test_support.py b/src/test/python/pylucene_test_support.py new file mode 100644 index 00000000..13494035 --- /dev/null +++ b/src/test/python/pylucene_test_support.py @@ -0,0 +1,904 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +from __future__ import annotations + +import os +import struct +import tempfile +from contextlib import ExitStack, contextmanager +from dataclasses import dataclass +from enum import Enum +from functools import lru_cache +from pathlib import Path +from typing import Any, Iterator + + +REPO_ROOT = Path(__file__).resolve().parents[3] +ID_FIELD = "id" +VECTOR_FIELD = "vector" +FILTER_FIELD = "filter-status" +FILTER_INCLUDED_VALUE = "included" +FILTER_EXCLUDED_VALUE = "excluded" + +CAGRA_TEST_CODEC_CLASS = ( + "com.nvidia.cuvs.lucene.PyLuceneTestSupport$CagraSearchCodec" +) +CAGRA_TEST_QUERY_CLASS = ( + "com.nvidia.cuvs.lucene.PyLuceneTestSupport$CagraSearchQuery" +) +HNSW_GRAPH_VERIFYING_QUERY_CLASS = ( + "com.nvidia.cuvs.lucene.PyLuceneTestSupport$HnswGraphVerifyingQuery" +) +CAGRA_VECTOR_READER_CLASS = "com.nvidia.cuvs.lucene.CuVS2510GPUVectorsReader" + +FORCE_CPU_HNSW_FALLBACK_PROPERTY = "cuvs.lucene.forceCpuHnswFallback" +QUERY_PROPERTIES = { + "field": "cuvs.lucene.pylucene.query.field", + "target": "cuvs.lucene.pylucene.query.target", + "k": "cuvs.lucene.pylucene.query.k", + "i_top_k": "cuvs.lucene.pylucene.query.iTopK", + "search_width": "cuvs.lucene.pylucene.query.searchWidth", + "expected_hnsw_m": "cuvs.lucene.pylucene.query.expectedHnswM", + "filter_field": "cuvs.lucene.pylucene.query.filterField", + "filter_value": "cuvs.lucene.pylucene.query.filterValue", +} + +_FLOAT32 = struct.Struct("!f") + + +@dataclass(frozen=True) +class IndexScenario: + name: str + codec_name: str + document_count: int + dimensions: int + top_k: int + segment_count: int = 1 + force_merge_segment_count: int = 0 + disable_automatic_merges: bool = True + document_ids_without_vectors: frozenset[int] = frozenset() + document_ids_to_delete: frozenset[int] = frozenset() + additional_query_document_ids: tuple[int, ...] = () + document_ids_accepted_by_filter: frozenset[int] = frozenset() + filter_query_document_id: int | None = None + force_cpu_hnsw: bool = False + codec_factory_class: str = "" + use_cagra_search_query: bool = False + expected_hnsw_m: int = 0 + search_width: int = 1 + i_top_k: int = 64 + + +class QueryDocumentState(Enum): + SEARCHABLE = "searchable" + WITHOUT_VECTOR = "without-vector" + DELETED = "deleted" + + +@dataclass(frozen=True) +class QueryObservation: + query_document_id: int + hit_ids: tuple[str, ...] + query_class: str + query_document_state: QueryDocumentState + + +@dataclass(frozen=True) +class FilteredQueryObservation: + query_document_id: int + hit_ids: tuple[str, ...] + query_class: str + + +@dataclass(frozen=True) +class VectorReaderMetadata: + vector_count: int + dimensions: tuple[int, ...] + reader_classes: tuple[str, ...] + hnsw_layer_counts: tuple[int, ...] + + +@dataclass(frozen=True) +class IndexRun: + index_files: tuple[str, ...] + pre_merge_segment_count: int + segment_count: int + live_document_count: int + max_document_count: int + vector_count: int + vector_dimensions: tuple[int, ...] + vector_reader_classes: tuple[str, ...] + hnsw_layer_counts: tuple[int, ...] + searchable_vector_document_ids: tuple[int, ...] + document_ids_without_vectors: frozenset[int] + deleted_document_ids: frozenset[int] + query_observations: tuple[QueryObservation, ...] + filtered_query_observation: FilteredQueryObservation | None + writer_telemetry: dict[str, str] + + +@dataclass +class PyLuceneContext: + cuvs_lucene_jar: Path + cuvs_java_jar: Path + test_classes: Path + codec_class: Any + query_class: Any + class_class: Any + system_class: Any + jarray: Any + codec_cache: dict[str, Any] + + +def find_cuvs_lucene_jar() -> Path: + configured = os.environ.get("CUVS_LUCENE_JAR") + if configured: + jar = Path(configured).resolve() + if not jar.is_file(): + raise FileNotFoundError(f"Configured cuvs-lucene jar does not exist: {jar}") + return jar + + jars = sorted( + jar + for jar in (REPO_ROOT / "target").glob("cuvs-lucene-*.jar") + if not any( + marker in jar.name + for marker in ("-jar-with-", "-sources.jar", "-javadoc.jar") + ) + ) + if not jars: + raise FileNotFoundError( + "No cuvs-lucene jar found under target/. " + "Run `mvn clean package -DskipTests` first." + ) + return jars[-1].resolve() + + +def find_cuvs_java_jar() -> Path: + configured = os.environ.get("CUVS_LUCENE_CUVS_JAVA_JAR") or os.environ.get( + "CUVS_JAVA_JAR" + ) + if configured: + jar = Path(configured).resolve() + if not jar.is_file(): + raise FileNotFoundError(f"Configured cuvs-java jar does not exist: {jar}") + return jar + + m2_repo = ( + Path.home() / ".m2" / "repository" / "com" / "nvidia" / "cuvs" / "cuvs-java" + ) + if not m2_repo.exists(): + raise FileNotFoundError( + "Unable to find cuvs-java in ~/.m2. Set " + "CUVS_LUCENE_CUVS_JAVA_JAR to the base cuvs-java jar." + ) + + def is_base_jar(jar: Path) -> bool: + return ( + jar.name.startswith("cuvs-java-") + and jar.name.endswith(".jar") + and "-x86_64-" not in jar.name + and "-sources" not in jar.name + and "-javadoc" not in jar.name + ) + + jars = sorted(jar for jar in m2_repo.glob("*/*.jar") if is_base_jar(jar)) + if not jars: + raise FileNotFoundError( + "Unable to find the base cuvs-java jar in ~/.m2. Set " + "CUVS_LUCENE_CUVS_JAVA_JAR explicitly." + ) + return jars[-1].resolve() + + +def find_test_classes() -> Path: + test_classes = Path( + os.environ.get( + "CUVS_LUCENE_PYLUCENE_TEST_CLASSES", + REPO_ROOT / "target" / "test-classes", + ) + ).resolve() + required_classes = ( + "PyLuceneTestSupport.class", + "PyLuceneTestSupport$CagraSearchCodec.class", + "PyLuceneTestSupport$CagraSearchQuery.class", + "PyLuceneTestSupport$HnswGraphVerifyingQuery.class", + ) + package_dir = test_classes / "com" / "nvidia" / "cuvs" / "lucene" + missing = [name for name in required_classes if not (package_dir / name).is_file()] + if missing: + raise FileNotFoundError( + f"Compiled PyLucene test bridge is incomplete under {test_classes}: " + f"missing {', '.join(missing)}. Run Maven test compilation first." + ) + return test_classes + + +@lru_cache(maxsize=None) +def deterministic_float32_vector( + document_id: int, dimensions: int +) -> tuple[float, ...]: + value = ((document_id + 1) * 2654435761) & 0xFFFFFFFF + vector = [] + for dimension in range(dimensions): + value = ( + 1664525 * value + 1013904223 + dimension * 17 + ) & 0xFFFFFFFF + component = (value / 4294967295.0) * 2.0 - 1.0 + vector.append(_FLOAT32.unpack(_FLOAT32.pack(component))[0]) + return tuple(vector) + + +def to_java_float_array(jarray: Any, values: tuple[float, ...]) -> Any: + return jarray("float")(values) + + +def _init_vm(cuvs_java_jar: Path, cuvs_lucene_jar: Path, test_classes: Path) -> Any: + import lucene + + java_library_path = os.environ.get("JAVA_LIBRARY_PATH") or os.environ.get( + "LD_LIBRARY_PATH" + ) + vmargs = [ + "--enable-native-access=ALL-UNNAMED", + "--add-modules=jdk.incubator.vector", + ( + "-Djava.util.logging.config.file=" + f"{REPO_ROOT / 'src' / 'main' / 'resources' / 'logging.properties'}" + ), + ] + if java_library_path: + vmargs.append(f"-Djava.library.path={java_library_path}") + + lucene.initVM( + classpath=os.pathsep.join( + [ + str(cuvs_java_jar), + str(cuvs_lucene_jar), + str(test_classes), + lucene.CLASSPATH, + ] + ), + vmargs=vmargs, + ) + return lucene + + +def initialize_pylucene_context( + expected_spi_codecs: tuple[str, ...], +) -> PyLuceneContext: + cuvs_lucene_jar = find_cuvs_lucene_jar() + cuvs_java_jar = find_cuvs_java_jar() + test_classes = find_test_classes() + lucene = _init_vm(cuvs_java_jar, cuvs_lucene_jar, test_classes) + + from java.lang import Class, System + from org.apache.lucene.codecs import Codec + from org.apache.lucene.search import Query + + # Resolve classpath failures before constructing codecs or indexes. + Class.forName("com.nvidia.cuvs.spi.JDKProvider") + Class.forName(CAGRA_TEST_CODEC_CLASS) + Class.forName(CAGRA_TEST_QUERY_CLASS) + Class.forName(HNSW_GRAPH_VERIFYING_QUERY_CLASS) + + available_codecs = Codec.availableCodecs() + for codec_name in expected_spi_codecs: + if not available_codecs.contains(codec_name): + raise RuntimeError( + f"{codec_name} was not advertised by Lucene SPI; " + f"available codecs: {available_codecs}" + ) + + return PyLuceneContext( + cuvs_lucene_jar=cuvs_lucene_jar, + cuvs_java_jar=cuvs_java_jar, + test_classes=test_classes, + codec_class=Codec, + query_class=Query, + class_class=Class, + system_class=System, + jarray=lucene.JArray, + codec_cache={}, + ) + + +def _codec_for_scenario(scenario: IndexScenario, context: PyLuceneContext) -> Any: + cache_key = scenario.codec_factory_class or f"spi:{scenario.codec_name}" + cached = context.codec_cache.get(cache_key) + if cached is not None: + return cached + + if scenario.codec_factory_class: + reflected = context.class_class.forName( + scenario.codec_factory_class + ).newInstance() + codec = context.codec_class.cast_(reflected) + else: + available_codecs = context.codec_class.availableCodecs() + if not available_codecs.contains(scenario.codec_name): + raise RuntimeError( + f"{scenario.codec_name} was not advertised by Lucene SPI; " + f"available codecs: {available_codecs}" + ) + codec = context.codec_class.forName(scenario.codec_name) + + if codec.getName() != scenario.codec_name: + raise RuntimeError( + f"Expected codec {scenario.codec_name}, got {codec.getName()}" + ) + context.codec_cache[cache_key] = codec + return codec + + +def _writer_telemetry(codec: Any) -> dict[str, str]: + description = str(codec.knnVectorsFormat()) + _, separator, payload = description.partition("(") + if not separator or not payload.endswith(")"): + raise RuntimeError(f"Malformed vector-format diagnostics: {description!r}") + + telemetry: dict[str, str] = {} + for item in payload[:-1].split(";"): + key, item_separator, value = item.partition("=") + if not item_separator or not key: + raise RuntimeError( + f"Malformed vector-format diagnostic item {item!r} in {description!r}" + ) + telemetry[key] = value + return telemetry + + +@contextmanager +def _temporary_system_properties( + system_class: Any, properties: dict[str, str] +) -> Iterator[None]: + previous = {name: system_class.getProperty(name) for name in properties} + try: + for name, value in properties.items(): + system_class.setProperty(name, value) + yield + finally: + for name, value in previous.items(): + if value is None: + system_class.clearProperty(name) + else: + system_class.setProperty(name, value) + + +@contextmanager +def _cpu_hnsw_fallback( + scenario: IndexScenario, system_class: Any +) -> Iterator[None]: + if not scenario.force_cpu_hnsw: + yield + return + + with _temporary_system_properties( + system_class, {FORCE_CPU_HNSW_FALLBACK_PROPERTY: "true"} + ): + yield + + +def _new_cagra_search_query( + context: PyLuceneContext, + target: tuple[float, ...], + top_k: int, + i_top_k: int, + search_width: int, + filter_field: str = "", + filter_value: str = "", +) -> Any: + properties = { + QUERY_PROPERTIES["field"]: VECTOR_FIELD, + QUERY_PROPERTIES["target"]: ",".join(f"{value:.9g}" for value in target), + QUERY_PROPERTIES["k"]: str(top_k), + QUERY_PROPERTIES["i_top_k"]: str(max(top_k, i_top_k)), + QUERY_PROPERTIES["search_width"]: str(search_width), + } + if filter_field: + properties[QUERY_PROPERTIES["filter_field"]] = filter_field + properties[QUERY_PROPERTIES["filter_value"]] = filter_value + with _temporary_system_properties(context.system_class, properties): + reflected = context.class_class.forName( + CAGRA_TEST_QUERY_CLASS + ).newInstance() + return context.query_class.cast_(reflected) + + +def _new_hnsw_graph_verifying_query( + context: PyLuceneContext, + target: tuple[float, ...], + top_k: int, + expected_m: int, + filter_field: str = "", + filter_value: str = "", +) -> Any: + properties = { + QUERY_PROPERTIES["field"]: VECTOR_FIELD, + QUERY_PROPERTIES["target"]: ",".join(f"{value:.9g}" for value in target), + QUERY_PROPERTIES["k"]: str(top_k), + QUERY_PROPERTIES["expected_hnsw_m"]: str(expected_m), + } + if filter_field: + properties[QUERY_PROPERTIES["filter_field"]] = filter_field + properties[QUERY_PROPERTIES["filter_value"]] = filter_value + with _temporary_system_properties(context.system_class, properties): + reflected = context.class_class.forName( + HNSW_GRAPH_VERIFYING_QUERY_CLASS + ).newInstance() + return context.query_class.cast_(reflected) + + +def _validate_document_ids(scenario: IndexScenario) -> None: + valid_document_ids = set(range(scenario.document_count)) + referenced_document_ids = ( + set(scenario.document_ids_without_vectors) + | set(scenario.document_ids_to_delete) + | set(scenario.additional_query_document_ids) + | set(scenario.document_ids_accepted_by_filter) + ) + if scenario.filter_query_document_id is not None: + referenced_document_ids.add(scenario.filter_query_document_id) + invalid_document_ids = referenced_document_ids - valid_document_ids + if invalid_document_ids: + raise ValueError( + f"{scenario.name}: document IDs are outside the index: " + f"{sorted(invalid_document_ids)}" + ) + + conflicting_document_ids = ( + scenario.document_ids_without_vectors & scenario.document_ids_to_delete + ) + if conflicting_document_ids: + raise ValueError( + f"{scenario.name}: documents cannot be both vectorless and deleted: " + f"{sorted(conflicting_document_ids)}" + ) + + if scenario.filter_query_document_id is None: + if scenario.document_ids_accepted_by_filter: + raise ValueError( + f"{scenario.name}: accepted filter documents require a " + "filtered query" + ) + return + + if not scenario.document_ids_accepted_by_filter: + raise ValueError( + f"{scenario.name}: filtered query has no accepted documents" + ) + if scenario.filter_query_document_id in ( + scenario.document_ids_without_vectors | scenario.document_ids_to_delete + ): + raise ValueError( + f"{scenario.name}: filtered query document must have a live " + f"vector: {scenario.filter_query_document_id}" + ) + if ( + scenario.filter_query_document_id + in scenario.document_ids_accepted_by_filter + ): + raise ValueError( + f"{scenario.name}: filtered query document must be rejected by " + "the filter" + ) + + +def _representative_query_document_ids( + scenario: IndexScenario, searchable_vector_document_ids: tuple[int, ...] +) -> tuple[int, ...]: + targets = (0, scenario.document_count // 2, scenario.document_count - 1) + query_document_ids: list[int] = [] + for target in targets: + nearest = min( + searchable_vector_document_ids, + key=lambda document_id: abs(document_id - target), + ) + if nearest not in query_document_ids: + query_document_ids.append(nearest) + return tuple(query_document_ids) + + +def _query_document_state( + query_document_id: int, + document_ids_without_vectors: frozenset[int], + deleted_document_ids: frozenset[int], +) -> QueryDocumentState: + if query_document_id in document_ids_without_vectors: + return QueryDocumentState.WITHOUT_VECTOR + if query_document_id in deleted_document_ids: + return QueryDocumentState.DELETED + return QueryDocumentState.SEARCHABLE + + +def segment_document_id_ranges( + document_count: int, segment_count: int +) -> tuple[range, ...]: + effective_segment_count = max(1, min(segment_count, document_count)) + documents_per_segment, remainder = divmod( + document_count, effective_segment_count + ) + ranges = [] + start_document_id = 0 + for segment_id in range(effective_segment_count): + documents_in_segment = documents_per_segment + ( + 1 if segment_id < remainder else 0 + ) + end_document_id = start_document_id + documents_in_segment + ranges.append(range(start_document_id, end_document_id)) + start_document_id = end_document_id + return tuple(ranges) + + +def _segment_end_document_ids(scenario: IndexScenario) -> set[int]: + return { + document_ids.stop - 1 + for document_ids in segment_document_id_ranges( + scenario.document_count, scenario.segment_count + ) + } + + +def _segment_count(directory: Any) -> int: + from org.apache.lucene.index import DirectoryReader + + reader = DirectoryReader.open(directory) + try: + return sum(1 for _ in reader.leaves()) + finally: + reader.close() + + +def _new_writer_config(codec: Any, suppress_merges: bool) -> Any: + from org.apache.lucene.index import ( + IndexWriterConfig, + LogDocMergePolicy, + NoMergePolicy, + SerialMergeScheduler, + ) + + config = IndexWriterConfig() + config.setCodec(codec) + config.setUseCompoundFile(False) + config.setMergeScheduler(SerialMergeScheduler()) + if suppress_merges: + config.setMergePolicy(NoMergePolicy.INSTANCE) + else: + merge_policy = LogDocMergePolicy() + merge_policy.setMergeFactor(1000) + config.setMergePolicy(merge_policy) + return config + + +def _hit_ids(stored_fields: Any, hits: Any) -> tuple[str, ...]: + return tuple(stored_fields.document(hit.doc).get(ID_FIELD) for hit in hits) + + +def _vector_reader_observations( + reader: Any, +) -> VectorReaderMetadata: + from org.apache.lucene.codecs.hnsw import HnswGraphProvider + from org.apache.lucene.index import SegmentReader + + vector_count = 0 + dimensions: list[int] = [] + reader_classes: list[str] = [] + hnsw_layer_counts: list[int] = [] + + for leaf_reader_context in reader.leaves(): + leaf_reader = leaf_reader_context.reader() + values = leaf_reader.getFloatVectorValues(VECTOR_FIELD) + if values is not None: + vector_count += values.size() + dimensions.append(values.dimension()) + + if not SegmentReader.instance_(leaf_reader): + reader_classes.append(str(leaf_reader.getClass().getName())) + continue + + segment_reader = SegmentReader.cast_(leaf_reader) + vector_reader = segment_reader.getVectorReader() + reader_classes.append(str(vector_reader.getClass().getName())) + if HnswGraphProvider.instance_(vector_reader): + graph_provider = HnswGraphProvider.cast_(vector_reader) + hnsw_layer_counts.append( + graph_provider.getGraph(VECTOR_FIELD).numLevels() + ) + + return VectorReaderMetadata( + vector_count=vector_count, + dimensions=tuple(dimensions), + reader_classes=tuple(reader_classes), + hnsw_layer_counts=tuple(hnsw_layer_counts), + ) + + +def _write_index_and_apply_deletions( + directory: Any, + codec: Any, + scenario: IndexScenario, + context: PyLuceneContext, + document_ids_without_vectors: frozenset[int], + document_ids_to_delete: frozenset[int], +) -> None: + from org.apache.lucene.document import ( + Document, + Field, + KnnFloatVectorField, + StringField, + ) + from org.apache.lucene.index import ( + IndexWriter, + Term, + VectorSimilarityFunction, + ) + + writer_config = _new_writer_config( + codec, suppress_merges=scenario.disable_automatic_merges + ) + writer = IndexWriter(directory, writer_config) + try: + segment_end_document_ids = _segment_end_document_ids(scenario) + for document_id in range(scenario.document_count): + document = Document() + document.add( + StringField(ID_FIELD, f"doc-{document_id}", Field.Store.YES) + ) + if document_id not in document_ids_without_vectors: + vector = deterministic_float32_vector(document_id, scenario.dimensions) + document.add( + KnnFloatVectorField( + VECTOR_FIELD, + to_java_float_array(context.jarray, vector), + VectorSimilarityFunction.EUCLIDEAN, + ) + ) + if scenario.filter_query_document_id is not None: + filter_value = ( + FILTER_INCLUDED_VALUE + if document_id + in scenario.document_ids_accepted_by_filter + else FILTER_EXCLUDED_VALUE + ) + document.add( + StringField( + FILTER_FIELD, + filter_value, + Field.Store.NO, + ) + ) + writer.addDocument(document) + if document_id in segment_end_document_ids: + writer.commit() + + for document_id in document_ids_to_delete: + writer.deleteDocuments(Term(ID_FIELD, f"doc-{document_id}")) + if document_ids_to_delete: + writer.commit() + finally: + writer.close() + + +def _force_merge_if_requested( + directory: Any, codec: Any, force_merge_segment_count: int +) -> None: + if not force_merge_segment_count: + return + + from org.apache.lucene.index import IndexWriter + + writer = IndexWriter( + directory, _new_writer_config(codec, suppress_merges=False) + ) + try: + writer.forceMerge(force_merge_segment_count) + writer.commit() + finally: + writer.close() + + +def _new_vector_query( + scenario: IndexScenario, + context: PyLuceneContext, + query_document_id: int, + top_k: int, + filter_field: str = "", + filter_value: str = "", +) -> Any: + from org.apache.lucene.index import Term + from org.apache.lucene.search import KnnFloatVectorQuery, TermQuery + + target = deterministic_float32_vector(query_document_id, scenario.dimensions) + if scenario.use_cagra_search_query: + return _new_cagra_search_query( + context, + target, + top_k, + scenario.i_top_k, + scenario.search_width, + filter_field, + filter_value, + ) + if scenario.expected_hnsw_m: + return _new_hnsw_graph_verifying_query( + context, + target, + top_k, + scenario.expected_hnsw_m, + filter_field, + filter_value, + ) + filter_query = ( + TermQuery(Term(filter_field, filter_value)) if filter_field else None + ) + if filter_query is None: + return KnnFloatVectorQuery( + VECTOR_FIELD, + to_java_float_array(context.jarray, target), + top_k, + ) + return KnnFloatVectorQuery( + VECTOR_FIELD, + to_java_float_array(context.jarray, target), + top_k, + filter_query, + ) + + +def _run_filtered_query( + searcher: Any, + stored_fields: Any, + context: PyLuceneContext, + scenario: IndexScenario, + query_document_id: int, + accepted_document_count: int, +) -> FilteredQueryObservation: + top_k = min(scenario.top_k, accepted_document_count) + query = _new_vector_query( + scenario, + context, + query_document_id, + top_k, + FILTER_FIELD, + FILTER_INCLUDED_VALUE, + ) + hits = searcher.search(query, top_k).scoreDocs + return FilteredQueryObservation( + query_document_id=query_document_id, + hit_ids=_hit_ids(stored_fields, hits), + query_class=str(query.getClass().getName()), + ) + + +def run_index_scenario( + scenario: IndexScenario, context: PyLuceneContext +) -> IndexRun: + from java.nio.file import Paths + from org.apache.lucene.index import DirectoryReader + from org.apache.lucene.search import IndexSearcher + from org.apache.lucene.store import FSDirectory + + codec = _codec_for_scenario(scenario, context) + _validate_document_ids(scenario) + document_ids_without_vectors = scenario.document_ids_without_vectors + document_ids_to_delete = scenario.document_ids_to_delete + searchable_vector_document_ids = tuple( + document_id + for document_id in range(scenario.document_count) + if document_id not in document_ids_without_vectors + and document_id not in document_ids_to_delete + ) + if not searchable_vector_document_ids: + raise RuntimeError( + f"{scenario.name}: no searchable vectors are available" + ) + filter_accepted_searchable_document_ids = tuple( + document_id + for document_id in searchable_vector_document_ids + if document_id in scenario.document_ids_accepted_by_filter + ) + if ( + scenario.filter_query_document_id is not None + and not filter_accepted_searchable_document_ids + ): + raise RuntimeError( + f"{scenario.name}: filter accepts no searchable vectors" + ) + representative_query_document_ids = _representative_query_document_ids( + scenario, searchable_vector_document_ids + ) + + with ExitStack() as stack: + stack.enter_context(_cpu_hnsw_fallback(scenario, context.system_class)) + index_path = stack.enter_context( + tempfile.TemporaryDirectory( + prefix=f"cuvs-lucene-pylucene-{scenario.name}-" + ) + ) + directory = FSDirectory.open(Paths.get(index_path)) + try: + _write_index_and_apply_deletions( + directory, + codec, + scenario, + context, + document_ids_without_vectors, + document_ids_to_delete, + ) + + pre_merge_segment_count = _segment_count(directory) + _force_merge_if_requested( + directory, codec, scenario.force_merge_segment_count + ) + + index_files = tuple( + sorted(path.name for path in Path(index_path).iterdir()) + ) + reader = DirectoryReader.open(directory) + try: + vector_reader_metadata = _vector_reader_observations(reader) + searcher = IndexSearcher(reader) + stored_fields = searcher.storedFields() + top_k = min(scenario.top_k, len(searchable_vector_document_ids)) + observations = [] + query_document_ids = tuple( + dict.fromkeys( + representative_query_document_ids + + scenario.additional_query_document_ids + ) + ) + for query_document_id in query_document_ids: + query = _new_vector_query( + scenario, context, query_document_id, top_k + ) + hits = searcher.search(query, top_k).scoreDocs + observations.append( + QueryObservation( + query_document_id=query_document_id, + hit_ids=_hit_ids(stored_fields, hits), + query_class=str(query.getClass().getName()), + query_document_state=_query_document_state( + query_document_id, + document_ids_without_vectors, + document_ids_to_delete, + ), + ) + ) + + filtered_query_observation = ( + _run_filtered_query( + searcher, + stored_fields, + context, + scenario, + scenario.filter_query_document_id, + len(filter_accepted_searchable_document_ids), + ) + if scenario.filter_query_document_id is not None + else None + ) + + result = IndexRun( + index_files=index_files, + pre_merge_segment_count=pre_merge_segment_count, + segment_count=sum(1 for _ in reader.leaves()), + live_document_count=reader.numDocs(), + max_document_count=reader.maxDoc(), + vector_count=vector_reader_metadata.vector_count, + vector_dimensions=vector_reader_metadata.dimensions, + vector_reader_classes=vector_reader_metadata.reader_classes, + hnsw_layer_counts=vector_reader_metadata.hnsw_layer_counts, + searchable_vector_document_ids=searchable_vector_document_ids, + document_ids_without_vectors=document_ids_without_vectors, + deleted_document_ids=document_ids_to_delete, + query_observations=tuple(observations), + filtered_query_observation=filtered_query_observation, + writer_telemetry=_writer_telemetry(codec), + ) + finally: + reader.close() + finally: + directory.close() + + return result diff --git a/src/test/python/test_pylucene_end_to_end.py b/src/test/python/test_pylucene_end_to_end.py new file mode 100644 index 00000000..dce063b5 --- /dev/null +++ b/src/test/python/test_pylucene_end_to_end.py @@ -0,0 +1,1245 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +from __future__ import annotations + +import os +from dataclasses import dataclass +from enum import Enum +from zipfile import ZipFile + +import pytest + +from pylucene_test_support import ( + CAGRA_TEST_CODEC_CLASS, + CAGRA_TEST_QUERY_CLASS, + CAGRA_VECTOR_READER_CLASS, + HNSW_GRAPH_VERIFYING_QUERY_CLASS, + IndexRun, + IndexScenario, + PyLuceneContext, + QueryDocumentState, + deterministic_float32_vector, + find_cuvs_lucene_jar, + initialize_pylucene_context, + run_index_scenario, + segment_document_id_ranges, +) + +# Python 3.14 reports one deprecation per JCC-generated PyLucene builtin type. +# This narrow third-party filter keeps native/cuVS warnings fully visible. +pytestmark = pytest.mark.filterwarnings( + "ignore:builtin type .* has no __module__ attribute:DeprecationWarning" +) + + +HNSW_CODEC = "Lucene101AcceleratedHNSWCodec" +CAGRA_HNSW_BASE_LAYER_CODEC = "Lucene101AcceleratedHNSWBaseLayerCodec" +CAGRA_HNSW_MULTI_LAYER_CODEC = "Lucene101AcceleratedHNSWMultiLayerCodec" +CAGRA_CODEC = "CuVS2510GPUSearchCodec" + +EXPECTED_SPI_CODECS = ( + HNSW_CODEC, + CAGRA_HNSW_BASE_LAYER_CODEC, + CAGRA_HNSW_MULTI_LAYER_CODEC, + CAGRA_CODEC, + "Lucene101AcceleratedHNSWBinaryQuantizedCodec", + "Lucene101AcceleratedHNSWScalarQuantizedCodec", +) +CODEC_SERVICE = "META-INF/services/org.apache.lucene.codecs.Codec" +VECTOR_FORMAT_SERVICE = ( + "META-INF/services/org.apache.lucene.codecs.KnnVectorsFormat" +) +EXPECTED_CODEC_PROVIDERS = frozenset( + { + "com.nvidia.cuvs.lucene.Lucene101AcceleratedHNSWCodec", + "com.nvidia.cuvs.lucene.Lucene101AcceleratedHNSWBaseLayerCodec", + "com.nvidia.cuvs.lucene.Lucene101AcceleratedHNSWMultiLayerCodec", + "com.nvidia.cuvs.lucene.CuVS2510GPUSearchCodec", + ( + "com.nvidia.cuvs.lucene." + "LuceneAcceleratedHNSWBinaryQuantizedCodec" + ), + ( + "com.nvidia.cuvs.lucene." + "LuceneAcceleratedHNSWScalarQuantizedCodec" + ), + } +) +EXPECTED_VECTOR_FORMAT_PROVIDERS = frozenset( + { + "com.nvidia.cuvs.lucene.CuVS2510GPUVectorsFormat", + "com.nvidia.cuvs.lucene.Lucene99AcceleratedHNSWVectorsFormat", + ( + "com.nvidia.cuvs.lucene." + "LuceneAcceleratedHNSWBinaryQuantizedVectorsFormat" + ), + ( + "com.nvidia.cuvs.lucene." + "LuceneAcceleratedHNSWScalarQuantizedVectorsFormat" + ), + } +) +REQUIRED_CUVS_LUCENE_CLASSES = frozenset( + provider.replace(".", "/") + ".class" + for provider in EXPECTED_CODEC_PROVIDERS | EXPECTED_VECTOR_FORMAT_PROVIDERS +) + +CAGRA_GRAPH_DEGREE = 32 +CAGRA_INTERMEDIATE_GRAPH_DEGREE = 64 +# cuVS expands NN-Descent's requested degree by 1.5 and reduces it with a +# warning when that internal degree is at least the number of input vectors. +NN_DESCENT_INTERNAL_GRAPH_DEGREE = CAGRA_INTERMEDIATE_GRAPH_DEGREE * 3 // 2 +MIN_VECTORS_PER_CAGRA_BUILD = NN_DESCENT_INTERNAL_GRAPH_DEGREE + 1 +CAGRA_HNSW_M = CAGRA_GRAPH_DEGREE // 2 +MIN_VECTORS_FOR_THREE_HNSW_LAYERS = ( + MIN_VECTORS_PER_CAGRA_BUILD * CAGRA_HNSW_M**2 +) +DEFAULT_MIN_RECALL = 0.75 + + +class ExecutionPath(Enum): + CPU_HNSW = ( + "CPU HNSW build -> HNSW search", + "cpu-hnsw", + "cpu-hnsw-fallback", + ) + GPU_CAGRA_BUILT_HNSW = ( + "GPU CAGRA build -> HNSW search", + "gpu-cagra-built-hnsw", + "gpu-hnsw", + ) + GPU_CAGRA_SEARCH = ( + "GPU CAGRA build -> CAGRA search", + "gpu-cagra-search", + "gpu-cagra", + ) + + def __init__( + self, label: str, selector: str, expected_writer_path: str + ) -> None: + self.label = label + self.selector = selector + self.expected_writer_path = expected_writer_path + + @property + def requires_gpu(self) -> bool: + return self is not ExecutionPath.CPU_HNSW + + +class DocumentSetup(Enum): + ALL_SEARCHABLE = "all-searchable" + ONE_DELETED = "one-deleted" + + +@dataclass(frozen=True) +class SuiteSettings: + requested_document_count: int + dimensions: int + top_k: int + minimum_recall: float + + +@dataclass(frozen=True) +class DocumentConfiguration: + document_ids_without_vectors: frozenset[int] = frozenset() + document_ids_to_delete: frozenset[int] = frozenset() + additional_query_document_ids: tuple[int, ...] = () + + +@dataclass(frozen=True) +class DocumentFilterConfiguration: + query_document_id: int | None = None + accepted_document_ids: frozenset[int] = frozenset() + + +@dataclass(frozen=True) +class EndToEndCase: + selector: str + pytest_id: str + selection_names: frozenset[str] + scenario: IndexScenario + execution_path: ExecutionPath + expected_index_file_suffixes: tuple[str, ...] + expected_hnsw_layers: int = 0 + min_recall: float = DEFAULT_MIN_RECALL + + +def _positive_int_env(name: str, default: int) -> int: + value = int(os.environ.get(name, default)) + if value <= 0: + raise ValueError(f"{name} must be positive, got {value}") + return value + + +def _boolean_env(name: str, default: bool) -> bool: + value = os.environ.get(name) + if value is None: + return default + normalized = value.strip().lower() + if normalized in {"1", "true", "yes", "on"}: + return True + if normalized in {"0", "false", "no", "off"}: + return False + raise ValueError(f"{name} must be a boolean value, got {value!r}") + + +def _minimum_recall_from_env() -> float: + value = float( + os.environ.get("CUVS_LUCENE_PYLUCENE_MIN_RECALL", DEFAULT_MIN_RECALL) + ) + if not 0.0 <= value <= 1.0: + raise ValueError( + "CUVS_LUCENE_PYLUCENE_MIN_RECALL must be between 0 and 1, " + f"got {value}" + ) + return value + + +def _suite_settings() -> SuiteSettings: + requested_document_count = _positive_int_env( + "CUVS_LUCENE_PYLUCENE_ROWS", 2000 + ) + dimensions = _positive_int_env("CUVS_LUCENE_PYLUCENE_DIMS", 32) + top_k = _positive_int_env("CUVS_LUCENE_PYLUCENE_TOPK", 20) + if dimensions > 4096: + raise ValueError( + f"CUVS_LUCENE_PYLUCENE_DIMS must be at most 4096, got {dimensions}" + ) + return SuiteSettings( + requested_document_count=requested_document_count, + dimensions=dimensions, + top_k=top_k, + minimum_recall=_minimum_recall_from_env(), + ) + + +def _document_configuration( + document_count: int, setup: DocumentSetup +) -> DocumentConfiguration: + middle_document_id = document_count // 2 + if setup is DocumentSetup.ONE_DELETED: + return DocumentConfiguration( + document_ids_to_delete=frozenset({middle_document_id}), + additional_query_document_ids=(middle_document_id,), + ) + return DocumentConfiguration() + + +def _minimum_document_count_for_selective_filter( + segment_count: int, top_k: int +) -> int: + return segment_count * 4 * (top_k + 1) + + +def _selective_document_filter_configuration( + document_count: int, segment_count: int, top_k: int +) -> DocumentFilterConfiguration: + accepted_document_ids = set() + for segment_document_ids in segment_document_id_ranges( + document_count, segment_count + ): + accepted_document_count = max( + top_k + 1, len(segment_document_ids) // 4 + ) + first_accepted_document_id = ( + segment_document_ids.stop - accepted_document_count + ) + accepted_document_ids.update( + range( + first_accepted_document_id, + segment_document_ids.stop, + ) + ) + + query_document_id = 0 + if query_document_id in accepted_document_ids: + raise ValueError( + "selective filter must reject its query document" + ) + return DocumentFilterConfiguration( + query_document_id=query_document_id, + accepted_document_ids=frozenset(accepted_document_ids), + ) + + +def _minimum_document_count( + execution_path: ExecutionPath, + segment_count: int, + hnsw_layers: int, + document_setup: DocumentSetup, +) -> int: + if execution_path is ExecutionPath.CPU_HNSW: + required_document_count = segment_count + elif hnsw_layers == 3: + required_document_count = MIN_VECTORS_FOR_THREE_HNSW_LAYERS + else: + required_document_count = ( + segment_count * MIN_VECTORS_PER_CAGRA_BUILD + ) + + if document_setup is DocumentSetup.ONE_DELETED: + required_document_count = max(required_document_count, 2) + return required_document_count + + +def _selection_names( + selector: str, + execution_path: ExecutionPath, + groups: tuple[str, ...], + legacy_aliases: tuple[str, ...], +) -> frozenset[str]: + names = { + selector, + execution_path.selector, + "all", + *groups, + *legacy_aliases, + } + if execution_path.requires_gpu: + names.add("gpu") + if execution_path is ExecutionPath.GPU_CAGRA_BUILT_HNSW: + names.add("gpu-cagra-hnsw") + names.add("cagra-hnsw") + return frozenset(names) + + +def _cpu_hnsw_case( + selector: str, + pytest_id: str, + *, + segment_count: int = 1, + force_merge_segment_count: int = 0, + document_setup: DocumentSetup = DocumentSetup.ALL_SEARCHABLE, + selective_filter: bool = False, + groups: tuple[str, ...] = (), + legacy_aliases: tuple[str, ...] = (), +) -> EndToEndCase: + settings = _suite_settings() + minimum_document_count = _minimum_document_count( + ExecutionPath.CPU_HNSW, + segment_count, + 0, + document_setup, + ) + if selective_filter: + minimum_document_count = max( + minimum_document_count, + _minimum_document_count_for_selective_filter( + segment_count, settings.top_k + ), + ) + document_count = max( + settings.requested_document_count, minimum_document_count + ) + documents = _document_configuration(document_count, document_setup) + document_filter = ( + _selective_document_filter_configuration( + document_count, segment_count, settings.top_k + ) + if selective_filter + else DocumentFilterConfiguration() + ) + scenario = IndexScenario( + name=selector, + codec_name=HNSW_CODEC, + document_count=document_count, + dimensions=settings.dimensions, + top_k=settings.top_k, + segment_count=segment_count, + force_merge_segment_count=force_merge_segment_count, + document_ids_without_vectors=documents.document_ids_without_vectors, + document_ids_to_delete=documents.document_ids_to_delete, + additional_query_document_ids=documents.additional_query_document_ids, + document_ids_accepted_by_filter=document_filter.accepted_document_ids, + filter_query_document_id=document_filter.query_document_id, + force_cpu_hnsw=True, + expected_hnsw_m=32, + ) + return EndToEndCase( + selector=selector, + pytest_id=pytest_id, + selection_names=_selection_names( + selector, + ExecutionPath.CPU_HNSW, + groups, + legacy_aliases, + ), + scenario=scenario, + execution_path=ExecutionPath.CPU_HNSW, + expected_index_file_suffixes=(".vex", ".vem"), + min_recall=settings.minimum_recall, + ) + + +def _cagra_built_hnsw_case( + selector: str, + pytest_id: str, + *, + hnsw_layers: int = 1, + segment_count: int = 1, + force_merge_segment_count: int = 0, + document_setup: DocumentSetup = DocumentSetup.ALL_SEARCHABLE, + selective_filter: bool = False, + groups: tuple[str, ...] = (), + legacy_aliases: tuple[str, ...] = (), +) -> EndToEndCase: + settings = _suite_settings() + minimum_document_count = _minimum_document_count( + ExecutionPath.GPU_CAGRA_BUILT_HNSW, + segment_count, + hnsw_layers, + document_setup, + ) + if selective_filter: + minimum_document_count = max( + minimum_document_count, + _minimum_document_count_for_selective_filter( + segment_count, settings.top_k + ), + ) + document_count = max( + settings.requested_document_count, minimum_document_count + ) + documents = _document_configuration(document_count, document_setup) + document_filter = ( + _selective_document_filter_configuration( + document_count, segment_count, settings.top_k + ) + if selective_filter + else DocumentFilterConfiguration() + ) + codec_name = ( + CAGRA_HNSW_MULTI_LAYER_CODEC + if hnsw_layers == 3 + else CAGRA_HNSW_BASE_LAYER_CODEC + ) + scenario = IndexScenario( + name=selector, + codec_name=codec_name, + document_count=document_count, + dimensions=settings.dimensions, + top_k=settings.top_k, + segment_count=segment_count, + force_merge_segment_count=force_merge_segment_count, + document_ids_without_vectors=documents.document_ids_without_vectors, + document_ids_to_delete=documents.document_ids_to_delete, + additional_query_document_ids=documents.additional_query_document_ids, + document_ids_accepted_by_filter=document_filter.accepted_document_ids, + filter_query_document_id=document_filter.query_document_id, + expected_hnsw_m=16, + ) + return EndToEndCase( + selector=selector, + pytest_id=pytest_id, + selection_names=_selection_names( + selector, + ExecutionPath.GPU_CAGRA_BUILT_HNSW, + groups, + legacy_aliases, + ), + scenario=scenario, + execution_path=ExecutionPath.GPU_CAGRA_BUILT_HNSW, + expected_index_file_suffixes=(".vex", ".vem"), + expected_hnsw_layers=hnsw_layers, + min_recall=settings.minimum_recall, + ) + + +def _cagra_search_case( + selector: str, + pytest_id: str, + *, + segment_count: int = 1, + force_merge_segment_count: int = 0, + document_setup: DocumentSetup = DocumentSetup.ALL_SEARCHABLE, + search_width: int = 1, + selective_filter: bool = False, + groups: tuple[str, ...] = (), + legacy_aliases: tuple[str, ...] = (), +) -> EndToEndCase: + settings = _suite_settings() + if settings.top_k > 1024: + raise ValueError( + "GPU CAGRA search cases require " + "CUVS_LUCENE_PYLUCENE_TOPK <= 1024" + ) + minimum_document_count = _minimum_document_count( + ExecutionPath.GPU_CAGRA_SEARCH, + segment_count, + 0, + document_setup, + ) + if selective_filter: + minimum_document_count = max( + minimum_document_count, + _minimum_document_count_for_selective_filter( + segment_count, settings.top_k + ), + ) + document_count = max( + settings.requested_document_count, minimum_document_count + ) + documents = _document_configuration(document_count, document_setup) + document_filter = ( + _selective_document_filter_configuration( + document_count, segment_count, settings.top_k + ) + if selective_filter + else DocumentFilterConfiguration() + ) + scenario = IndexScenario( + name=selector, + codec_name=CAGRA_CODEC, + codec_factory_class=CAGRA_TEST_CODEC_CLASS, + document_count=document_count, + dimensions=settings.dimensions, + top_k=settings.top_k, + segment_count=segment_count, + force_merge_segment_count=force_merge_segment_count, + document_ids_without_vectors=documents.document_ids_without_vectors, + document_ids_to_delete=documents.document_ids_to_delete, + additional_query_document_ids=documents.additional_query_document_ids, + document_ids_accepted_by_filter=document_filter.accepted_document_ids, + filter_query_document_id=document_filter.query_document_id, + use_cagra_search_query=True, + search_width=search_width, + i_top_k=max(64, settings.top_k), + ) + return EndToEndCase( + selector=selector, + pytest_id=pytest_id, + selection_names=_selection_names( + selector, + ExecutionPath.GPU_CAGRA_SEARCH, + groups, + legacy_aliases, + ), + scenario=scenario, + execution_path=ExecutionPath.GPU_CAGRA_SEARCH, + expected_index_file_suffixes=(".vcag", ".vemc"), + min_recall=settings.minimum_recall, + ) + + +SEGMENT_CASES = ( + _cpu_hnsw_case( + "cpu-hnsw-1-segment", + "cpu-hnsw-1-segment", + groups=( + "execution-paths", + "algorithm-matrix", + "segment-topologies", + ), + legacy_aliases=("smoke", "hnsw-cpu", "hnsw-cpu-1seg"), + ), + _cpu_hnsw_case( + "cpu-hnsw-10-segments", + "cpu-hnsw-10-segments", + segment_count=10, + groups=("segment-topologies",), + legacy_aliases=("hnsw-cpu-10seg",), + ), + _cagra_built_hnsw_case( + "gpu-cagra-built-hnsw-1-segment", + "gpu-cagra-built-hnsw-1-segment", + groups=( + "execution-paths", + "algorithm-matrix", + "segment-topologies", + "hnsw-layer-counts", + ), + legacy_aliases=( + "hnsw-1seg", + "cagra-hnsw-1layer", + "cagra-hnsw-base", + "cagra-hnsw-base-layer", + ), + ), + _cagra_built_hnsw_case( + "gpu-cagra-built-hnsw-10-segments", + "gpu-cagra-built-hnsw-10-segments", + segment_count=10, + groups=("segment-topologies",), + legacy_aliases=("hnsw-10seg",), + ), + _cagra_search_case( + "gpu-cagra-search-10-segments", + "gpu-cagra-search-10-segments", + segment_count=10, + groups=("segment-topologies",), + legacy_aliases=("cagra-10seg",), + ), +) + +FORCE_MERGE_CASES = ( + _cpu_hnsw_case( + "cpu-hnsw-10-to-1-force-merge", + "cpu-hnsw-10-to-1", + segment_count=10, + force_merge_segment_count=1, + groups=("force-merges", "segment-topologies"), + legacy_aliases=("hnsw-cpu-10seg-force-1",), + ), + _cpu_hnsw_case( + "cpu-hnsw-100-to-10-force-merge", + "cpu-hnsw-100-to-10", + segment_count=100, + force_merge_segment_count=10, + groups=("force-merges", "segment-topologies"), + legacy_aliases=("hnsw-cpu-100seg-force-10",), + ), + _cagra_built_hnsw_case( + "gpu-cagra-built-hnsw-10-to-1-force-merge", + "gpu-cagra-built-hnsw-10-to-1", + segment_count=10, + force_merge_segment_count=1, + groups=("force-merges", "segment-topologies"), + legacy_aliases=("hnsw-10seg-force-1",), + ), + _cagra_built_hnsw_case( + "gpu-cagra-built-hnsw-100-to-10-force-merge", + "gpu-cagra-built-hnsw-100-to-10", + segment_count=100, + force_merge_segment_count=10, + groups=("force-merges", "segment-topologies"), + legacy_aliases=("hnsw-100seg-force-10",), + ), + _cagra_search_case( + "gpu-cagra-search-10-to-1-force-merge", + "gpu-cagra-search-10-to-1", + segment_count=10, + force_merge_segment_count=1, + groups=("force-merges", "segment-topologies"), + legacy_aliases=("cagra-10seg-force-1",), + ), + _cagra_search_case( + "gpu-cagra-search-100-to-10-force-merge", + "gpu-cagra-search-100-to-10", + segment_count=100, + force_merge_segment_count=10, + groups=("force-merges", "segment-topologies"), + legacy_aliases=("cagra-100seg-force-10",), + ), +) + +HNSW_LAYER_CASES = ( + _cagra_built_hnsw_case( + "gpu-cagra-built-hnsw-3-layers", + "3-layers", + hnsw_layers=3, + groups=("hnsw-layer-counts",), + legacy_aliases=( + "cagra-hnsw-3layer", + "cagra-hnsw-multilayer", + "cagra-hnsw-multi-layer", + ), + ), +) + +CAGRA_SEARCH_WIDTH_CASES = ( + _cagra_search_case( + "gpu-cagra-search-width-1", + "1", + groups=( + "execution-paths", + "algorithm-matrix", + "segment-topologies", + "cagra-search-widths", + ), + legacy_aliases=( + "gpu-cagra-search-1-segment", + "cagra-1seg", + ), + ), + _cagra_search_case( + "gpu-cagra-search-width-16", + "16", + search_width=16, + groups=("cagra-search-widths",), + ), + _cagra_search_case( + "gpu-cagra-search-width-32", + "32", + search_width=32, + groups=("cagra-search-widths",), + ), +) + +DELETED_DOCUMENT_CASES = ( + _cagra_search_case( + "gpu-cagra-search-deleted-documents", + "gpu-cagra-search", + document_setup=DocumentSetup.ONE_DELETED, + groups=("deleted-documents",), + ), +) + +DOCUMENT_FILTER_CASES = ( + _cpu_hnsw_case( + "cpu-hnsw-selective-filter", + "cpu-hnsw-selective-filter", + selective_filter=True, + groups=("document-filter",), + legacy_aliases=("cpu-hnsw-document-filter",), + ), + _cagra_built_hnsw_case( + "gpu-cagra-built-hnsw-selective-filter", + "gpu-cagra-built-hnsw-selective-filter", + selective_filter=True, + groups=("document-filter",), + legacy_aliases=("gpu-cagra-built-hnsw-document-filter",), + ), + _cagra_search_case( + "gpu-cagra-search-selective-filter-10-segments", + "gpu-cagra-search-selective-filter-10-segments", + segment_count=10, + selective_filter=True, + groups=("document-filter",), + legacy_aliases=("gpu-cagra-search-document-filter",), + ), +) + + +def _case_parameter(case: EndToEndCase) -> object: + marker = pytest.mark.pylucene_case(*sorted(case.selection_names)) + return pytest.param(case, id=case.pytest_id, marks=marker) + + +def _case_parameters( + cases: tuple[EndToEndCase, ...], +) -> tuple[object, ...]: + return tuple(_case_parameter(case) for case in cases) + + +def _service_providers( + archive: ZipFile, service_path: str +) -> frozenset[str]: + lines = archive.read(service_path).decode("utf-8").splitlines() + return frozenset( + line.strip() + for line in lines + if line.strip() and not line.lstrip().startswith("#") + ) + + +def test_published_jar_has_expected_lucene_services() -> None: + cuvs_lucene_jar = find_cuvs_lucene_jar() + with ZipFile(cuvs_lucene_jar) as archive: + entries = frozenset(archive.namelist()) + expected_services = frozenset( + {CODEC_SERVICE, VECTOR_FORMAT_SERVICE} + ) + assert expected_services <= entries, ( + f"{cuvs_lucene_jar.name} is missing service descriptors: " + f"{sorted(expected_services - entries)}" + ) + assert REQUIRED_CUVS_LUCENE_CLASSES <= entries, ( + f"{cuvs_lucene_jar.name} is missing classes: " + f"{sorted(REQUIRED_CUVS_LUCENE_CLASSES - entries)}" + ) + + lucene_service_descriptors = { + entry + for entry in entries + if entry.startswith("META-INF/services/org.apache.lucene.") + } + assert lucene_service_descriptors == expected_services + + codec_providers = _service_providers(archive, CODEC_SERVICE) + vector_format_providers = _service_providers( + archive, VECTOR_FORMAT_SERVICE + ) + assert EXPECTED_CODEC_PROVIDERS <= codec_providers + assert EXPECTED_VECTOR_FORMAT_PROVIDERS <= vector_format_providers + assert not any( + provider.startswith("org.apache.lucene.") + for provider in codec_providers | vector_format_providers + ) + + +def test_published_jar_does_not_bundle_lucene_or_cuvs_java() -> None: + cuvs_lucene_jar = find_cuvs_lucene_jar() + with ZipFile(cuvs_lucene_jar) as archive: + entries = frozenset(archive.namelist()) + + bundled_lucene_classes = sorted( + entry + for entry in entries + if entry.startswith("org/apache/lucene/") and not entry.endswith("/") + ) + assert not bundled_lucene_classes, ( + "PyLucene must supply Lucene classes; the cuvs-lucene jar contains " + f"{bundled_lucene_classes}" + ) + + bundled_cuvs_java_classes = sorted( + entry + for entry in entries + if entry.startswith("com/nvidia/cuvs/") + and not entry.startswith("com/nvidia/cuvs/lucene/") + and not entry.endswith("/") + ) + assert not bundled_cuvs_java_classes, ( + "The base cuvs-java jar must remain separate; cuvs-lucene contains " + f"{bundled_cuvs_java_classes}" + ) + + bundled_multi_release_classes = sorted( + entry + for entry in entries + if entry.startswith("META-INF/versions/") + and "/com/nvidia/cuvs/" in entry + and not entry.endswith("/") + ) + assert not bundled_multi_release_classes, ( + "The base multi-release cuvs-java jar must remain separate; " + f"cuvs-lucene contains {bundled_multi_release_classes}" + ) + + +def test_published_jar_excludes_pylucene_test_support() -> None: + cuvs_lucene_jar = find_cuvs_lucene_jar() + with ZipFile(cuvs_lucene_jar) as archive: + entries = frozenset(archive.namelist()) + + test_support_prefix = "com/nvidia/cuvs/lucene/PyLuceneTestSupport" + published_test_classes = sorted( + entry for entry in entries if entry.startswith(test_support_prefix) + ) + assert not published_test_classes + + +@pytest.fixture(scope="session") +def pylucene_context() -> PyLuceneContext: + expected_codecs = ( + EXPECTED_SPI_CODECS + if _boolean_env("CUVS_LUCENE_VERIFY_ALL_CODECS", True) + else () + ) + return initialize_pylucene_context(expected_codecs) + + +def _squared_distance( + left: tuple[float, ...], right: tuple[float, ...] +) -> float: + return sum( + (left_value - right_value) ** 2 + for left_value, right_value in zip(left, right) + ) + + +def _brute_force_neighbor_ids( + case: EndToEndCase, + result: IndexRun, + query_document_id: int, + candidate_document_ids: tuple[int, ...] | None = None, +) -> tuple[str, ...]: + query_vector = deterministic_float32_vector( + query_document_id, case.scenario.dimensions + ) + candidates = ( + result.searchable_vector_document_ids + if candidate_document_ids is None + else candidate_document_ids + ) + ordered_ids = sorted( + candidates, + key=lambda document_id: ( + _squared_distance( + query_vector, + deterministic_float32_vector( + document_id, case.scenario.dimensions + ), + ), + document_id, + ), + ) + return tuple( + f"doc-{document_id}" + for document_id in ordered_ids[: case.scenario.top_k] + ) + + +def _assert_index_files_and_segments( + case: EndToEndCase, result: IndexRun +) -> None: + for suffix in case.expected_index_file_suffixes: + assert any(name.endswith(suffix) for name in result.index_files), ( + f"{case.selector}: no index file ending with {suffix}; " + f"files={result.index_files}" + ) + + expected_initial_segments = min( + case.scenario.segment_count, case.scenario.document_count + ) + assert result.pre_merge_segment_count == expected_initial_segments, ( + f"{case.selector}: expected {expected_initial_segments} segments " + f"before forceMerge, got {result.pre_merge_segment_count}" + ) + expected_final_segments = ( + case.scenario.force_merge_segment_count or expected_initial_segments + ) + assert result.segment_count == expected_final_segments, ( + f"{case.selector}: expected {expected_final_segments} final segments, " + f"got {result.segment_count}" + ) + + +def _assert_index_metadata(case: EndToEndCase, result: IndexRun) -> None: + expected_live_documents = ( + case.scenario.document_count - len(result.deleted_document_ids) + ) + assert result.live_document_count == expected_live_documents + expected_max_documents = ( + expected_live_documents + if case.scenario.force_merge_segment_count + else case.scenario.document_count + ) + assert result.max_document_count == expected_max_documents + + expected_vector_count = ( + len(result.searchable_vector_document_ids) + if case.scenario.force_merge_segment_count + else ( + case.scenario.document_count + - len(result.document_ids_without_vectors) + ) + ) + assert result.vector_count == expected_vector_count, ( + f"{case.selector}: expected {expected_vector_count} vector values, " + f"got {result.vector_count}" + ) + assert result.vector_dimensions + assert set(result.vector_dimensions) == {case.scenario.dimensions} + + +def _assert_execution_path(case: EndToEndCase, result: IndexRun) -> None: + observed_writer_path = result.writer_telemetry.get("writerPath") + assert observed_writer_path == case.execution_path.expected_writer_path + + if case.execution_path is ExecutionPath.GPU_CAGRA_SEARCH: + assert set(result.vector_reader_classes) == {CAGRA_VECTOR_READER_CLASS} + assert all( + observation.query_class == CAGRA_TEST_QUERY_CLASS + for observation in result.query_observations + ) + if result.filtered_query_observation is not None: + assert ( + result.filtered_query_observation.query_class + == CAGRA_TEST_QUERY_CLASS + ) + assert not result.hnsw_layer_counts + return + + assert result.vector_reader_classes + assert all( + observation.query_class == HNSW_GRAPH_VERIFYING_QUERY_CLASS + for observation in result.query_observations + ) + if result.filtered_query_observation is not None: + assert ( + result.filtered_query_observation.query_class + == HNSW_GRAPH_VERIFYING_QUERY_CLASS + ) + assert all( + "Lucene99HnswVectorsReader" in reader_class + for reader_class in result.vector_reader_classes + ) + assert all( + not reader_class.startswith("com.nvidia.cuvs.lucene") + for reader_class in result.vector_reader_classes + ) + + +def _assert_graph_configuration( + case: EndToEndCase, result: IndexRun +) -> None: + if case.expected_hnsw_layers: + assert result.hnsw_layer_counts + assert set(result.hnsw_layer_counts) == { + case.expected_hnsw_layers + }, ( + f"{case.selector}: expected {case.expected_hnsw_layers} HNSW " + f"layers, got {result.hnsw_layer_counts}" + ) + + if not case.execution_path.requires_gpu: + return + + telemetry = result.writer_telemetry + assert telemetry.get("cagraGraphBuildAlgo") == "NN_DESCENT" + assert telemetry.get("cagraGraphDegree") == "32" + assert telemetry.get("cagraIntermediateGraphDegree") == "64" + if case.execution_path is ExecutionPath.GPU_CAGRA_SEARCH: + assert telemetry.get("cagraStrategy") == "CUSTOM" + + +def _assert_search_results( + case: EndToEndCase, result: IndexRun +) -> tuple[float, ...]: + documents_without_vectors = { + f"doc-{document_id}" + for document_id in result.document_ids_without_vectors + } + deleted_documents = { + f"doc-{document_id}" for document_id in result.deleted_document_ids + } + expected_hit_count = min( + case.scenario.top_k, len(result.searchable_vector_document_ids) + ) + + recalls = [] + for observation in result.query_observations: + hits = observation.hit_ids + queried_document = f"doc-{observation.query_document_id}" + if ( + observation.query_document_state + is QueryDocumentState.SEARCHABLE + ): + assert len(hits) == expected_hit_count, ( + f"{case.selector}: query {queried_document} expected " + f"{expected_hit_count} hits, got {hits}" + ) + assert hits[0] == queried_document, ( + f"{case.selector}: queried document must rank first; " + f"expected {queried_document}, got {hits}" + ) + else: + assert queried_document not in hits, ( + f"{case.selector}: {observation.query_document_state.value} " + f"document was returned: {hits}" + ) + + assert len(hits) == len(set(hits)), ( + f"{case.selector}: duplicate hits returned: {hits}" + ) + returned_vectorless_documents = tuple( + document + for document in hits + if document in documents_without_vectors + ) + assert not returned_vectorless_documents, ( + f"{case.selector}: live documents without vectors were returned: " + f"{returned_vectorless_documents}" + ) + returned_deleted_documents = tuple( + document for document in hits if document in deleted_documents + ) + assert not returned_deleted_documents, ( + f"{case.selector}: deleted documents were returned: " + f"{returned_deleted_documents}" + ) + + if ( + observation.query_document_state + is not QueryDocumentState.SEARCHABLE + ): + continue + + expected_neighbors = _brute_force_neighbor_ids( + case, result, observation.query_document_id + ) + recall = ( + len(set(hits) & set(expected_neighbors)) + / len(expected_neighbors) + ) + assert recall >= case.min_recall, ( + f"{case.selector}: query {queried_document} recall " + f"{recall:.3f} is below floor {case.min_recall:.3f}; " + f"expected={expected_neighbors}, actual={hits}" + ) + recalls.append(recall) + + return tuple(recalls) + + +def _print_result( + case: EndToEndCase, result: IndexRun, recalls: tuple[float, ...] +) -> None: + report_label = case.selector.removeprefix( + f"{case.execution_path.selector}-" + ) + segment_summary = ( + f"{result.pre_merge_segment_count}->{result.segment_count}" + if case.scenario.force_merge_segment_count + else str(result.segment_count) + ) + details = [ + f"documents={case.scenario.document_count}", + f"liveDocuments={result.live_document_count}", + f"segments={segment_summary}", + f"topK={case.scenario.top_k}", + f"recall(min/avg)={min(recalls):.3f}/" + f"{sum(recalls) / len(recalls):.3f}", + ] + if result.document_ids_without_vectors: + details.append( + f"documentsWithoutVectors={len(result.document_ids_without_vectors)}" + ) + if result.deleted_document_ids: + details.append(f"deletedDocuments={len(result.deleted_document_ids)}") + if case.execution_path is ExecutionPath.GPU_CAGRA_SEARCH: + details.append(f"searchWidth={case.scenario.search_width}") + else: + details.append( + f"hnswLayers={sorted(set(result.hnsw_layer_counts))}" + ) + details.append(f"hnswM={case.scenario.expected_hnsw_m}") + + print( + f"PASS [{case.execution_path.label}] {report_label}: " + + ", ".join(details) + ) + + +def _run_and_verify( + case: EndToEndCase, context: PyLuceneContext +) -> tuple[IndexRun, tuple[float, ...]]: + result = run_index_scenario(case.scenario, context) + _assert_index_files_and_segments(case, result) + _assert_index_metadata(case, result) + _assert_execution_path(case, result) + _assert_graph_configuration(case, result) + recalls = _assert_search_results(case, result) + _print_result(case, result, recalls) + return result, recalls + + +@pytest.mark.parametrize("case", _case_parameters(SEGMENT_CASES)) +def test_search_with_configured_segment_count( + pylucene_context: PyLuceneContext, case: EndToEndCase +) -> None: + result, _ = _run_and_verify(case, pylucene_context) + assert not result.document_ids_without_vectors + assert not result.deleted_document_ids + + +@pytest.mark.parametrize("case", _case_parameters(FORCE_MERGE_CASES)) +def test_search_after_force_merge( + pylucene_context: PyLuceneContext, case: EndToEndCase +) -> None: + result, _ = _run_and_verify(case, pylucene_context) + assert not result.document_ids_without_vectors + assert not result.deleted_document_ids + + +@pytest.mark.parametrize("case", _case_parameters(HNSW_LAYER_CASES)) +def test_cagra_built_hnsw_has_expected_layer_count( + pylucene_context: PyLuceneContext, case: EndToEndCase +) -> None: + result, _ = _run_and_verify(case, pylucene_context) + assert set(result.hnsw_layer_counts) == {case.expected_hnsw_layers} + + +@pytest.mark.parametrize( + "case", _case_parameters(CAGRA_SEARCH_WIDTH_CASES) +) +def test_cagra_search_with_configured_search_width( + pylucene_context: PyLuceneContext, case: EndToEndCase +) -> None: + _run_and_verify(case, pylucene_context) + + +@pytest.mark.parametrize( + "case", _case_parameters(DELETED_DOCUMENT_CASES) +) +def test_deleted_documents_are_not_searchable( + pylucene_context: PyLuceneContext, case: EndToEndCase +) -> None: + result, _ = _run_and_verify(case, pylucene_context) + assert len(result.deleted_document_ids) == 1 + assert not result.document_ids_without_vectors + assert any( + observation.query_document_state is QueryDocumentState.DELETED + for observation in result.query_observations + ) + + +@pytest.mark.parametrize( + "case", _case_parameters(DOCUMENT_FILTER_CASES) +) +def test_vector_search_honors_selective_document_filter( + pylucene_context: PyLuceneContext, case: EndToEndCase +) -> None: + result, _ = _run_and_verify(case, pylucene_context) + observation = result.filtered_query_observation + assert observation is not None + + accepted_document_ids = tuple( + document_id + for document_id in result.searchable_vector_document_ids + if document_id in case.scenario.document_ids_accepted_by_filter + ) + accepted_document_id_set = set(accepted_document_ids) + accepted_hit_ids = { + f"doc-{document_id}" for document_id in accepted_document_ids + } + accepted_counts_by_segment = tuple( + sum( + document_id in accepted_document_id_set + for document_id in segment_document_ids + ) + for segment_document_ids in segment_document_id_ranges( + case.scenario.document_count, + case.scenario.segment_count, + ) + ) + query_document_id = case.scenario.filter_query_document_id + assert query_document_id is not None + queried_document = f"doc-{query_document_id}" + + assert observation.query_document_id == query_document_id + assert query_document_id not in ( + case.scenario.document_ids_accepted_by_filter + ) + assert queried_document not in observation.hit_ids + assert min(accepted_counts_by_segment) > case.scenario.top_k, ( + f"{case.selector}: every segment must retain more than topK=" + f"{case.scenario.top_k} accepted vectors; " + f"acceptedPerSegment={accepted_counts_by_segment}" + ) + + rejected_hit_ids = tuple( + hit_id + for hit_id in observation.hit_ids + if hit_id not in accepted_hit_ids + ) + assert not rejected_hit_ids, ( + f"{case.selector}: filter-rejected documents were returned: " + f"{rejected_hit_ids}" + ) + expected_hit_count = min( + case.scenario.top_k, + len(accepted_document_ids), + ) + assert len(observation.hit_ids) == expected_hit_count, ( + f"{case.selector}: expected {expected_hit_count} filtered hits, " + f"got {observation.hit_ids}" + ) + assert len(observation.hit_ids) == len(set(observation.hit_ids)), ( + f"{case.selector}: duplicate filtered hits returned: " + f"{observation.hit_ids}" + ) + + expected_neighbors = _brute_force_neighbor_ids( + case, + result, + query_document_id, + accepted_document_ids, + ) + recall = ( + len(set(observation.hit_ids) & set(expected_neighbors)) + / len(expected_neighbors) + ) + assert recall >= case.min_recall, ( + f"{case.selector}: filtered query {queried_document} recall " + f"{recall:.3f} is below floor {case.min_recall:.3f}; " + f"expected={expected_neighbors}, actual={observation.hit_ids}" + ) + print( + f"FILTER [{case.execution_path.label}] {case.selector}: " + f"acceptedPerSegment={min(accepted_counts_by_segment)}-" + f"{max(accepted_counts_by_segment)}, " + f"recall={recall:.3f}" + ) diff --git a/test_pylucene.sh b/test_pylucene.sh new file mode 100755 index 00000000..f9968295 --- /dev/null +++ b/test_pylucene.sh @@ -0,0 +1,9 @@ +#!/bin/bash + +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +set -euo pipefail + +REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +exec "${REPO_ROOT}/ci/run_pylucene_pytests.sh" "$@"