From 269909b61c912bf22b50c3084dd980921873055c Mon Sep 17 00:00:00 2001 From: nvzm123 Date: Tue, 6 Oct 2026 04:25:32 +0000 Subject: [PATCH 1/2] Add Lucene benchmark backend core --- .../cuvs_bench/cuvs_bench/backends/lucene.py | 1638 +++++++++++++++++ .../config/algos/lucene_accelerated_hnsw.yaml | 13 + .../config/algos/lucene_cpu_hnsw.yaml | 13 + .../config/algos/lucene_cuvs_cagra.yaml | 13 + .../cuvs_bench/tests/_lucene_test_support.py | 407 ++++ .../tests/test_lucene_backend_contract.py | 72 + .../tests/test_lucene_build_search.py | 837 +++++++++ .../cuvs_bench/tests/test_lucene_lifecycle.py | 708 +++++++ .../tests/test_lucene_results_config.py | 293 +++ 9 files changed, 3994 insertions(+) create mode 100644 python/cuvs_bench/cuvs_bench/backends/lucene.py create mode 100644 python/cuvs_bench/cuvs_bench/config/algos/lucene_accelerated_hnsw.yaml create mode 100644 python/cuvs_bench/cuvs_bench/config/algos/lucene_cpu_hnsw.yaml create mode 100644 python/cuvs_bench/cuvs_bench/config/algos/lucene_cuvs_cagra.yaml create mode 100644 python/cuvs_bench/cuvs_bench/tests/_lucene_test_support.py create mode 100644 python/cuvs_bench/cuvs_bench/tests/test_lucene_backend_contract.py create mode 100644 python/cuvs_bench/cuvs_bench/tests/test_lucene_build_search.py create mode 100644 python/cuvs_bench/cuvs_bench/tests/test_lucene_lifecycle.py create mode 100644 python/cuvs_bench/cuvs_bench/tests/test_lucene_results_config.py diff --git a/python/cuvs_bench/cuvs_bench/backends/lucene.py b/python/cuvs_bench/cuvs_bench/backends/lucene.py new file mode 100644 index 0000000000..77d728b429 --- /dev/null +++ b/python/cuvs_bench/cuvs_bench/backends/lucene.py @@ -0,0 +1,1638 @@ +# +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# + +"""Opt-in Lucene backend using stock PyLucene and cuvs-lucene codecs.""" + +from __future__ import annotations + +import hashlib +import json +import math +import os +import re +import shutil +import tempfile +import time +import uuid +from dataclasses import dataclass +from numbers import Integral +from pathlib import Path +from typing import Any, Callable, Mapping, Optional, Sequence + +import numpy as np + +from cuvs_bench._bin_format import read_bin_header + +from ..orchestrator.config_loaders import ( + BenchmarkConfig, + ConfigLoader, + IndexConfig, +) +from ._lucene_runtime import ( + ACCELERATED_HNSW_CODEC, + CAGRA_CODEC, + CPU_HNSW_CODEC, + DIRECT_PYLUCENE_DISPATCH, + MAX_CAGRA_TOP_K, + LuceneRuntime, + RuntimeBuildResult, + RuntimeSearchResult, + RuntimeSearchTiming, + TIMED_BRIDGE_PYLUCENE_DISPATCH, +) +from ._lucene_runtime_config import resolve_lucene_runtime_config +from ._utils import dtype_from_filename +from .base import BenchmarkBackend, BuildResult, Dataset, SearchResult + +CPU_HNSW_ALGORITHM = "lucene_cpu_hnsw" +ACCELERATED_HNSW_ALGORITHM = "lucene_accelerated_hnsw" +CAGRA_ALGORITHM = "lucene_cuvs_cagra" + +MAX_CPU_HNSW_DIMENSIONS = 1024 +MAX_CAGRA_DIMENSIONS = 4096 + +_CODEC_BY_ALGORITHM = { + CPU_HNSW_ALGORITHM: CPU_HNSW_CODEC, + ACCELERATED_HNSW_ALGORITHM: ACCELERATED_HNSW_CODEC, + CAGRA_ALGORITHM: CAGRA_CODEC, +} +_MAX_DIMENSIONS_BY_ALGORITHM = { + CPU_HNSW_ALGORITHM: MAX_CPU_HNSW_DIMENSIONS, + ACCELERATED_HNSW_ALGORITHM: MAX_CAGRA_DIMENSIONS, + CAGRA_ALGORITHM: MAX_CAGRA_DIMENSIONS, +} +_CUVS_ALGORITHMS = frozenset((ACCELERATED_HNSW_ALGORITHM, CAGRA_ALGORITHM)) +_BUILD_ROUTE_POLICY_BY_ALGORITHM = { + CPU_HNSW_ALGORITHM: "cpu_hnsw", + ACCELERATED_HNSW_ALGORITHM: "gpu_cagra_or_cpu_hnsw_fallback", + CAGRA_ALGORITHM: "gpu_cagra", +} +_SEARCH_ROUTE_BY_ALGORITHM = { + CPU_HNSW_ALGORITHM: "cpu_hnsw", + ACCELERATED_HNSW_ALGORITHM: "cpu_hnsw", + CAGRA_ALGORITHM: "gpu_cagra", +} +_MANIFEST_FILE = ".cuvs-bench-lucene.json" +_MANIFEST_SCHEMA = 2 +_SAFE_LABEL = re.compile(r"[A-Za-z0-9][A-Za-z0-9_.-]*") +_RUNTIME_KEYS = ( + "cuvs_java_jar", + "cuvs_lucene_jar", + "java_library_path", + "jvm_args", +) +_SCORE_ROUNDOFF_TOLERANCE = float(np.spacing(np.float32(1.0))) +_INDEX_PREWARM_BLOCK_BYTES = 1024 * 1024 +_TIMING_CONTRACT_VERSION = 1 + + +@dataclass(frozen=True) +class _IndexPrewarmTiming: + wall_ns: int + bytes_read: int + file_count: int + + +def _nanoseconds_to_milliseconds(value: int) -> float: + return float(value) / 1_000_000.0 + + +def _nanoseconds_to_seconds(value: int) -> float: + return float(value) / 1_000_000_000.0 + + +def _timing_statistics( + scope: str, samples_ns: list[int] +) -> dict[str, float | int]: + """Summarize one timing sample per measured query.""" + if not samples_ns or any(value < 0 for value in samples_ns): + raise RuntimeError(f"Lucene returned invalid {scope} timing data") + samples_ms = np.asarray(samples_ns, dtype=np.float64) / 1_000_000.0 + return { + f"{scope}_count": len(samples_ns), + f"{scope}_total_ms": float(samples_ms.sum()), + f"{scope}_mean_ms": float(samples_ms.mean()), + f"{scope}_p50_ms": float(np.percentile(samples_ms, 50)), + f"{scope}_p95_ms": float(np.percentile(samples_ms, 95)), + f"{scope}_p99_ms": float(np.percentile(samples_ms, 99)), + } + + +def _validate_nested_timing( + parent_scope: str, + parent_ns: int, + child_timings_ns: Mapping[str, int], +) -> None: + """Reject impossible same-clock timing relationships.""" + if parent_ns < 0 or any(value < 0 for value in child_timings_ns.values()): + raise RuntimeError( + f"Lucene returned invalid {parent_scope} timing data" + ) + child_total_ns = sum(child_timings_ns.values()) + if child_total_ns > parent_ns: + raise RuntimeError( + f"Lucene returned impossible {parent_scope} timing data: " + f"child phases total {child_total_ns} ns, exceeding " + f"the enclosing {parent_ns} ns" + ) + + +def _prewarm_index_files(index_path: Path) -> _IndexPrewarmTiming: + """Populate the host page cache by sequentially reading index files.""" + started = time.perf_counter_ns() + bytes_read = 0 + file_count = 0 + buffer = bytearray(_INDEX_PREWARM_BLOCK_BYTES) + for path in sorted(index_path.iterdir()): + if path.is_symlink(): + raise RuntimeError( + f"Lucene index prewarm refuses symbolic links: {path}" + ) + if not path.is_file(): + continue + file_count += 1 + try: + with path.open("rb", buffering=0) as stream: + while count := stream.readinto(buffer): + bytes_read += count + except OSError as error: + raise RuntimeError( + f"Failed to prewarm Lucene index file: {path}" + ) from error + if file_count == 0: + raise RuntimeError( + f"Lucene index contains no files to prewarm: {index_path}" + ) + return _IndexPrewarmTiming( + wall_ns=time.perf_counter_ns() - started, + bytes_read=bytes_read, + file_count=file_count, + ) + + +def _runtime_build_timing_metadata( + result: RuntimeBuildResult, +) -> dict[str, float]: + timing = result.timing + phases = { + "runtime_directory_open_seconds": timing.directory_open_ns, + "runtime_writer_setup_seconds": timing.writer_setup_ns, + "runtime_document_ingest_seconds": timing.document_ingest_ns, + "runtime_writer_commit_close_seconds": timing.writer_commit_close_ns, + "runtime_post_build_reader_seconds": timing.post_build_reader_ns, + "runtime_directory_close_seconds": timing.directory_close_ns, + } + _validate_nested_timing( + "runtime_build_wall", + timing.runtime_build_wall_ns, + phases, + ) + fields = { + **phases, + "runtime_build_wall_seconds": timing.runtime_build_wall_ns, + } + return { + name: _nanoseconds_to_seconds(value) for name, value in fields.items() + } + + +def _java_search_timing_metadata( + timing: RuntimeSearchTiming, +) -> dict[str, Any]: + """Validate the dispatch route and summarize independent JVM timings.""" + java_samples = [ + item.java_index_searcher_search_ns for item in timing.queries + ] + java_available = all(value is not None for value in java_samples) + if any(value is not None for value in java_samples) and not java_available: + raise RuntimeError( + "Lucene returned Java timing data for only some queries" + ) + valid_dispatch_kinds = { + DIRECT_PYLUCENE_DISPATCH, + TIMED_BRIDGE_PYLUCENE_DISPATCH, + } + if timing.search_dispatch_kind not in valid_dispatch_kinds: + raise RuntimeError( + "Lucene returned an unknown search dispatch kind: " + f"{timing.search_dispatch_kind!r}" + ) + expected_java_timing = ( + timing.search_dispatch_kind == TIMED_BRIDGE_PYLUCENE_DISPATCH + ) + if java_available != expected_java_timing: + raise RuntimeError( + "Lucene search dispatch kind does not match Java timing availability" + ) + + metadata: dict[str, Any] = { + "search_dispatch_kind": timing.search_dispatch_kind, + "java_timing_available": java_available, + } + if not java_available: + return metadata + + resolved_java_samples = [int(value) for value in java_samples] + metadata.update( + _timing_statistics( + "java_index_searcher_search", + resolved_java_samples, + ) + ) + metadata.update( + { + "first_query_java_index_searcher_search_ms": ( + _nanoseconds_to_milliseconds(resolved_java_samples[0]) + ), + "subsequent_java_index_searcher_search_mean_ms": ( + float(np.mean(resolved_java_samples[1:])) / 1_000_000.0 + if len(resolved_java_samples) > 1 + else None + ), + } + ) + return metadata + + +def _runtime_search_timing_metadata( + result: RuntimeSearchResult, expected_query_count: int +) -> tuple[dict[str, Any], list[int]]: + timing = result.timing + if len(timing.queries) != expected_query_count: + raise RuntimeError( + "Lucene returned timing data for " + f"{len(timing.queries)} queries, expected {expected_query_count}" + ) + plan_phases = { + "directory_open_ms": timing.directory_open_ns, + "reader_searcher_setup_ms": timing.reader_searcher_setup_ns, + "query_corpus_wall_ms": timing.query_corpus_wall_ns, + "reader_close_ms": timing.reader_close_ns, + "directory_close_ms": timing.directory_close_ns, + } + _validate_nested_timing( + "runtime_plan_wall", + timing.runtime_plan_wall_ns, + plan_phases, + ) + if timing.query_corpus_wall_ns <= 0: + raise RuntimeError("Lucene returned an empty query-corpus duration") + + for query_number, query_timing in enumerate(timing.queries): + _validate_nested_timing( + f"client_query[{query_number}]", + query_timing.client_query_ns, + { + "query_prepare": query_timing.query_prepare_ns, + "pylucene_search_dispatch": ( + query_timing.pylucene_search_dispatch_ns + ), + "result_materialization": ( + query_timing.result_materialization_ns + ), + }, + ) + client_query_total_ns = sum( + item.client_query_ns for item in timing.queries + ) + if client_query_total_ns > timing.query_corpus_wall_ns: + raise RuntimeError( + "Lucene returned impossible query_corpus_wall timing data: " + f"client queries total {client_query_total_ns} ns, exceeding " + f"the enclosing {timing.query_corpus_wall_ns} ns" + ) + + scopes = { + "query_prepare": [item.query_prepare_ns for item in timing.queries], + "pylucene_search_dispatch": [ + item.pylucene_search_dispatch_ns for item in timing.queries + ], + "result_materialization": [ + item.result_materialization_ns for item in timing.queries + ], + "client_query": [item.client_query_ns for item in timing.queries], + } + plan_timings = { + **plan_phases, + "runtime_plan_wall_ms": timing.runtime_plan_wall_ns, + } + metadata: dict[str, Any] = { + name: _nanoseconds_to_milliseconds(value) + for name, value in plan_timings.items() + } + for scope, samples in scopes.items(): + metadata.update(_timing_statistics(scope, samples)) + metadata.update(_java_search_timing_metadata(timing)) + + dispatch_samples = scopes["pylucene_search_dispatch"] + client_samples = scopes["client_query"] + metadata.update( + { + "first_query_pylucene_search_dispatch_ms": ( + _nanoseconds_to_milliseconds(dispatch_samples[0]) + ), + "first_query_client_query_ms": _nanoseconds_to_milliseconds( + client_samples[0] + ), + "subsequent_query_count": max(0, len(dispatch_samples) - 1), + "subsequent_pylucene_search_dispatch_mean_ms": ( + float(np.mean(dispatch_samples[1:])) / 1_000_000.0 + if len(dispatch_samples) > 1 + else None + ), + "subsequent_client_query_mean_ms": ( + float(np.mean(client_samples[1:])) / 1_000_000.0 + if len(client_samples) > 1 + else None + ), + } + ) + return metadata, client_samples + + +def _validate_vectors( + vectors: Any, label: str, *, maximum_dimensions: int +) -> np.ndarray: + array = np.asarray(vectors) + if array.ndim != 2 or not array.shape[0] or not array.shape[1]: + raise ValueError(f"{label} must be a nonempty two-dimensional array") + if array.dtype != np.float32: + raise TypeError(f"{label} must use float32 values, got {array.dtype}") + if array.shape[1] > maximum_dimensions: + raise ValueError( + f"{label} dimensions must not exceed {maximum_dimensions}" + ) + if not np.isfinite(array).all(): + raise ValueError(f"{label} must contain only finite values") + return np.ascontiguousarray(array) + + +def _validate_metric(dataset: Dataset) -> None: + if dataset.distance_metric.casefold() not in {"euclidean", "l2"}: + raise ValueError( + "The Lucene backend currently supports only Euclidean/L2 data, " + f"not {dataset.distance_metric!r}" + ) + + +def _score_to_squared_euclidean(score: float) -> float: + """Invert Lucene's ``score = 1 / (1 + squared_distance)`` transform.""" + if ( + not math.isfinite(score) + or score <= 0.0 + or score > 1.0 + _SCORE_ROUNDOFF_TOLERANCE + ): + raise RuntimeError( + f"Lucene returned an invalid Euclidean score: {score}" + ) + # Admit one float32 ULP defensively at the mathematical upper boundary. + score = min(score, 1.0) + return max(0.0, (1.0 / score) - 1.0) + + +def _format_exception(error: BaseException) -> str: + details = [f"{type(error).__name__}: {error}"] + details.extend(getattr(error, "__notes__", ())) + return "\n".join(details) + + +def _safe_label(value: str, kind: str) -> str: + if not isinstance(value, str) or _SAFE_LABEL.fullmatch(value) is None: + raise ValueError(f"Unsafe Lucene {kind}: {value!r}") + return value + + +def _codec_for(algorithm: str, build_params: Mapping[str, Any]) -> str: + try: + expected = _CODEC_BY_ALGORITHM[algorithm] + except KeyError as error: + raise ValueError( + f"Unsupported Lucene algorithm: {algorithm!r}" + ) from error + unsupported = set(build_params) - {"codec"} + if unsupported: + raise ValueError( + "Unsupported Lucene build parameters: " + + ", ".join(sorted(str(name) for name in unsupported)) + ) + actual = build_params.get("codec", expected) + if actual != expected: + raise ValueError( + f"{algorithm} requires codec {expected!r}, got {actual!r}" + ) + return expected + + +def _search_parameters( + algorithm: str, parameters: Mapping[str, Any], k: int +) -> dict[str, int]: + if not isinstance(parameters, Mapping): + raise TypeError("Lucene search parameters must be mappings") + if algorithm == CAGRA_ALGORITHM: + if parameters: + raise ValueError( + "lucene_cuvs_cagra currently uses the codec's fixed search " + "defaults and accepts no search parameters" + ) + if k > MAX_CAGRA_TOP_K: + raise ValueError( + f"CAGRA search supports k <= {MAX_CAGRA_TOP_K}; k={k} would " + "use the codec's brute-force route" + ) + return {"num_candidates": k} + unsupported = set(parameters) - {"num_candidates"} + if unsupported: + raise ValueError( + "Unsupported Lucene search parameters: " + + ", ".join(sorted(str(name) for name in unsupported)) + ) + candidates = parameters.get("num_candidates", k) + if type(candidates) is not int or candidates < k: + raise ValueError( + f"num_candidates must be an integer >= k ({k}), got {candidates!r}" + ) + return {"num_candidates": candidates} + + +def _source_identity(dataset: Dataset) -> dict[str, Any]: + source = Path(dataset.base_file).expanduser().resolve() + stat = source.stat() + return { + "path": str(source), + "device": stat.st_dev, + "size": stat.st_size, + "mtime_ns": stat.st_mtime_ns, + "ctime_ns": stat.st_ctime_ns, + "inode": stat.st_ino, + } + + +def _normalize_subset_size(dataset: Dataset) -> int | None: + value = dataset.metadata.get("subset_size") + if value is None: + return None + if isinstance(value, bool) or not isinstance(value, Integral) or value < 1: + raise ValueError( + f"subset_size must be a positive integer, got {value!r}" + ) + normalized = int(value) + if type(value) is not int: + dataset.metadata = {**dataset.metadata, "subset_size": normalized} + return normalized + + +def _dataset_identity(dataset: Dataset, vectors: np.ndarray) -> dict[str, Any]: + identity: dict[str, Any] = { + "name": dataset.name, + "vector_count": int(vectors.shape[0]), + "dimensions": int(vectors.shape[1]), + "subset_size": _normalize_subset_size(dataset), + "sha256": hashlib.sha256(vectors.view(np.uint8)).hexdigest(), + } + if dataset.base_file: + identity["source"] = _source_identity(dataset) + return identity + + +def _file_backed_dataset_identity( + dataset: Dataset, + stored: Mapping[str, Any], + *, + maximum_dimensions: int, +) -> tuple[dict[str, Any], int, int]: + """Stream a file source identity without materializing its vectors.""" + subset_size = _normalize_subset_size(dataset) + source = _source_identity(dataset) + dtype = np.dtype(dtype_from_filename(source["path"])) + if dtype != np.dtype(np.float32): + raise TypeError( + f"training vectors must use float32 values, got {dtype}" + ) + rows, dimensions, header_bytes = read_bin_header( + source["path"], dtype.itemsize + ) + if dimensions > maximum_dimensions: + raise ValueError( + f"training vectors dimensions must not exceed {maximum_dimensions}" + ) + if subset_size is not None: + rows = min(rows, subset_size) + remaining = int(rows) * int(dimensions) * dtype.itemsize + digest_builder = hashlib.sha256() + with Path(source["path"]).open("rb") as stream: + stream.seek(header_bytes) + while remaining: + chunk = stream.read(min(remaining, 8 * 1024 * 1024)) + if not chunk: + raise ValueError( + f"training vector file is truncated: {source['path']}" + ) + digest_builder.update(chunk) + remaining -= len(chunk) + confirmed_source = _source_identity(dataset) + if confirmed_source != source: + raise RuntimeError( + "training vector file changed while its identity was being read" + ) + digest = digest_builder.hexdigest() + stored_digest = stored.get("sha256") + if ( + not isinstance(stored_digest, str) + or re.fullmatch(r"[0-9a-f]{64}", stored_digest) is None + ): + raise RuntimeError( + "Lucene index manifest has an invalid dataset digest" + ) + identity = { + "name": dataset.name, + "vector_count": int(rows), + "dimensions": int(dimensions), + "subset_size": subset_size, + "sha256": digest, + "source": source, + } + return identity, int(rows), int(dimensions) + + +def _manifest_payload( + dataset: Dataset, + vectors: np.ndarray, + algorithm: str, + codec: str, + segment_count: int, + build_runtime_artifacts: Mapping[str, str], +) -> dict[str, Any]: + return { + "schema_version": _MANIFEST_SCHEMA, + "algorithm": algorithm, + "codec": codec, + "dataset": _dataset_identity(dataset, vectors), + "segment_count": segment_count, + "build_runtime_artifacts": dict(build_runtime_artifacts), + } + + +def _artifact_metadata( + role: str, provenance: Mapping[str, str] +) -> dict[str, str]: + return {f"{role}_{key}": value for key, value in provenance.items()} + + +def _write_manifest(index_path: Path, payload: Mapping[str, Any]) -> None: + target = index_path / _MANIFEST_FILE + temporary = index_path / f"{_MANIFEST_FILE}.tmp" + temporary.write_text( + json.dumps(payload, indent=2, sort_keys=True) + "\n", encoding="utf-8" + ) + temporary.chmod(0o644) + os.replace(temporary, target) + + +def _is_manifest_integer(value: Any, *, minimum: int = 0) -> bool: + return type(value) is int and value >= minimum + + +def _validate_manifest_dataset(dataset: Any, path: Path) -> None: + required = { + "name", + "vector_count", + "dimensions", + "subset_size", + "sha256", + } + if not isinstance(dataset, dict) or set(dataset) not in ( + required, + required | {"source"}, + ): + raise RuntimeError( + f"Lucene index manifest has an invalid dataset: {path}" + ) + if not isinstance(dataset["name"], str): + raise RuntimeError( + f"Lucene index manifest has an invalid dataset: {path}" + ) + for field in ("vector_count", "dimensions"): + if not _is_manifest_integer(dataset[field], minimum=1): + raise RuntimeError( + f"Lucene index manifest has an invalid dataset {field}: {path}" + ) + subset_size = dataset["subset_size"] + if subset_size is not None and not _is_manifest_integer( + subset_size, minimum=1 + ): + raise RuntimeError( + f"Lucene index manifest has an invalid dataset subset_size: {path}" + ) + digest = dataset["sha256"] + if ( + not isinstance(digest, str) + or re.fullmatch(r"[0-9a-f]{64}", digest) is None + ): + raise RuntimeError( + f"Lucene index manifest has an invalid dataset digest: {path}" + ) + if "source" not in dataset: + return + source = dataset["source"] + source_fields = { + "path", + "device", + "size", + "mtime_ns", + "ctime_ns", + "inode", + } + if not isinstance(source, dict) or set(source) != source_fields: + raise RuntimeError( + f"Lucene index manifest has an invalid dataset source: {path}" + ) + if ( + not isinstance(source["path"], str) + or not Path(source["path"]).is_absolute() + ): + raise RuntimeError( + f"Lucene index manifest has an invalid dataset source path: {path}" + ) + for field in source_fields - {"path"}: + if not _is_manifest_integer(source[field]): + raise RuntimeError( + "Lucene index manifest has an invalid dataset source " + f"{field}: {path}" + ) + + +def _read_manifest(index_path: Path) -> dict[str, Any]: + path = index_path / _MANIFEST_FILE + try: + payload = json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError) as error: + raise RuntimeError( + f"Cannot read Lucene index manifest {path}: {error}" + ) from error + if ( + not isinstance(payload, dict) + or type(payload.get("schema_version")) is not int + or payload["schema_version"] != _MANIFEST_SCHEMA + ): + raise RuntimeError( + f"Unsupported Lucene index manifest: {path}; rerun with --force" + ) + required = { + "schema_version", + "algorithm", + "codec", + "dataset", + "segment_count", + "build_runtime_artifacts", + } + if set(payload) != required: + raise RuntimeError(f"Lucene index manifest has invalid fields: {path}") + if not isinstance(payload["algorithm"], str) or not isinstance( + payload["codec"], str + ): + raise RuntimeError( + f"Lucene index manifest has invalid identifiers: {path}" + ) + _validate_manifest_dataset(payload["dataset"], path) + segment_count = payload.get("segment_count") + if not _is_manifest_integer(segment_count, minimum=1): + raise RuntimeError( + f"Lucene index manifest has an invalid segment count: {path}" + ) + if not isinstance(payload["build_runtime_artifacts"], dict): + raise RuntimeError( + "Lucene index manifest has invalid build-runtime artifact " + f"provenance: {path}" + ) + return payload + + +def _index_size(index_path: Path) -> int: + return sum( + path.stat().st_size for path in index_path.rglob("*") if path.is_file() + ) + + +class LuceneConfigLoader(ConfigLoader): + """Load the fixed initial Lucene algorithms from standard Bench YAML.""" + + def __init__(self, config_path: Optional[str] = None): + self.config_path = config_path or str( + Path(__file__).resolve().parents[1] / "config" + ) + + @property + def backend_type(self) -> str: + return "lucene" + + @staticmethod + def _split(value: Any) -> list[str] | None: + if not value: + return None + return [part.strip() for part in str(value).split(",") if part.strip()] + + def _discover_algo_groups( + self, + dataset_conf: dict, + dataset: str, + dataset_path: str, + **kwargs: Any, + ) -> list[tuple[str, str, dict, dict]]: + files = self.gather_algorithm_configs( + self.config_path, kwargs.get("algorithm_configuration") + ) + configs = {} + for path in files: + config = self.load_yaml_file(path) + if ( + isinstance(config, dict) + and config.get("name") in _CODEC_BY_ALGORITHM + ): + configs[config["name"]] = config + + algorithms = self._split(kwargs.get("algorithms")) + groups = self._split(kwargs.get("groups")) or ["base"] + requested_pairs: dict[str, list[str]] = {} + if algo_groups := self._split(kwargs.get("algo_groups")): + for item in algo_groups: + algorithm, separator, group = item.partition(".") + if not separator: + raise ValueError( + f"Lucene --algo-groups entry must be algorithm.group: {item!r}" + ) + selected = requested_pairs.setdefault(algorithm, []) + if group not in selected: + selected.append(group) + + requested_algorithms = ( + algorithms or list(requested_pairs) or list(configs) + ) + result = [] + for algorithm in requested_algorithms: + if algorithm not in _CODEC_BY_ALGORITHM: + raise ValueError( + f"Unsupported Lucene algorithm: {algorithm!r}" + ) + config = configs.get(algorithm) + if config is None: + raise ValueError(f"No configuration found for {algorithm!r}") + selected_groups = requested_pairs.get(algorithm, groups) + for group in selected_groups: + try: + group_config = config["groups"][group] + except KeyError as error: + raise ValueError( + f"No Lucene group {group!r} for {algorithm!r}" + ) from error + result.append((algorithm, group, group_config, {})) + return result + + def _build_benchmark_configs( + self, + dataset_config: Any, + dataset_conf: dict, + dataset: str, + dataset_path: str, + expanded_groups: list[tuple], + **kwargs: Any, + ) -> list[BenchmarkConfig]: + include_cuvs = any( + item[0] in _CUVS_ALGORITHMS for item in expanded_groups + ) + runtime_overrides = { + key: kwargs[key] + for key in _RUNTIME_KEYS + if kwargs.get(key) is not None + } + _safe_label(dataset, "dataset") + dataset_base = Path(dataset_path).expanduser().resolve() + root = (dataset_base / dataset / "index").resolve() + if not root.is_relative_to(dataset_base): + raise ValueError( + f"Lucene index root escapes the dataset path: {root}" + ) + configs = [] + for ( + algorithm, + group, + _raw, + build_combos, + search_combos, + _meta, + ) in expanded_groups: + _safe_label(algorithm, "algorithm") + _safe_label(group, "group") + for build_params in build_combos: + codec = _codec_for(algorithm, build_params) + name = algorithm if group == "base" else f"{algorithm}_{group}" + index = IndexConfig( + name=name, + algo=algorithm, + build_param={"codec": codec}, + search_params=[dict(item) for item in search_combos], + file=str(root / name), + ) + configs.append( + BenchmarkConfig( + indexes=[index], + backend_config={ + "name": name, + "algo": algorithm, + "group": group, + "codec": codec, + "index_root": str(root), + "requires_cuvs": algorithm in _CUVS_ALGORITHMS, + "include_cuvs": include_cuvs, + **runtime_overrides, + }, + ) + ) + return configs + + +class LuceneBackend(BenchmarkBackend): + """Build and query one immutable Lucene vector index per configuration.""" + + default_algorithm = CAGRA_ALGORITHM + + def __init__( + self, + config: dict[str, Any], + runtime_factory: Callable[ + [Mapping[str, Any]], LuceneRuntime + ] = LuceneRuntime.create, + ): + super().__init__(config) + self.algorithm = str(config.get("algo", "")) + self.codec = str(config.get("codec", "")) + self.group = _safe_label(str(config.get("group", "base")), "group") + if _CODEC_BY_ALGORITHM.get(self.algorithm) != self.codec: + raise ValueError( + f"Invalid Lucene algorithm/codec pair: {self.algorithm!r}, {self.codec!r}" + ) + self.maximum_dimensions = _MAX_DIMENSIONS_BY_ALGORITHM[self.algorithm] + self._runtime_factory = runtime_factory + self._runtime: LuceneRuntime | None = None + + @classmethod + def result_failure_message( + cls, results: Sequence[BuildResult | SearchResult] + ) -> str | None: + """Fail the public command when any Lucene operation failed.""" + failures = [result for result in results if not result.success] + if not failures: + return None + return "; ".join( + result.error_message or "unknown Lucene backend failure" + for result in failures + ) + + @property + def algo(self) -> str: + """Return the selected cuVS Bench algorithm name.""" + return self.algorithm + + def _get_runtime(self) -> LuceneRuntime: + if self._runtime is None: + resolved = resolve_lucene_runtime_config( + self.config, + requires_cuvs=bool(self.config.get("requires_cuvs")), + include_cuvs=bool(self.config.get("include_cuvs")), + ) + self._runtime = self._runtime_factory(resolved) + return self._runtime + + def _index(self, indexes: list[IndexConfig]) -> IndexConfig: + if len(indexes) != 1: + raise ValueError( + "The Lucene backend expects one index configuration" + ) + index = indexes[0] + if index.algo != self.algorithm: + raise ValueError( + f"Index algorithm {index.algo!r} does not match {self.algorithm!r}" + ) + _codec_for(index.algo, index.build_param) + return index + + def _index_path(self, index: IndexConfig) -> Path: + path = Path(os.path.abspath(Path(index.file).expanduser())) + root = Path(self.config["index_root"]).expanduser().resolve() + if path.is_symlink(): + raise ValueError( + f"Lucene index path must not be a symlink: {path}" + ) + if path.parent.resolve() != root or path == root: + raise ValueError( + f"Lucene index path is outside its configured root: {path}" + ) + return path + + def _failure_build( + self, index: IndexConfig, error: Exception + ) -> BuildResult: + return BuildResult( + index_path=index.file, + build_time_seconds=0.0, + index_size_bytes=0, + algorithm=index.algo, + build_params=dict(index.build_param), + success=False, + error_message=_format_exception(error), + metadata={"group": self.group, "index_name": index.name}, + ) + + def _validate_manifest_identity( + self, payload: Mapping[str, Any], dataset_identity: Mapping[str, Any] + ) -> None: + build_runtime_artifacts = payload.get("build_runtime_artifacts") + if not isinstance(build_runtime_artifacts, Mapping): + raise RuntimeError( + "Lucene index manifest has invalid build-runtime artifact provenance" + ) + artifact_keys = { + "cuvs_java_coordinates", + "cuvs_java_jar_path", + "cuvs_java_jar_sha256", + "cuvs_lucene_coordinates", + "cuvs_lucene_jar_path", + "cuvs_lucene_jar_sha256", + } + if ( + build_runtime_artifacts + and set(build_runtime_artifacts) != artifact_keys + ): + raise RuntimeError( + "Lucene index manifest has incomplete build-runtime artifact provenance" + ) + if self.algorithm in _CUVS_ALGORITHMS and not build_runtime_artifacts: + raise RuntimeError( + "cuVS-backed Lucene index manifest has no build-runtime " + "artifact provenance" + ) + if build_runtime_artifacts: + expected_coordinates = { + "cuvs_java_coordinates": "com.nvidia.cuvs:cuvs-java:", + "cuvs_lucene_coordinates": ( + "com.nvidia.cuvs.lucene:cuvs-lucene:" + ), + } + for key, prefix in expected_coordinates.items(): + value = build_runtime_artifacts[key] + if ( + not isinstance(value, str) + or not value.startswith(prefix) + or len(value) == len(prefix) + ): + raise RuntimeError( + f"Lucene index manifest has invalid {key}" + ) + for key in ("cuvs_java_jar_path", "cuvs_lucene_jar_path"): + value = build_runtime_artifacts[key] + if not isinstance(value, str) or not Path(value).is_absolute(): + raise RuntimeError( + f"Lucene index manifest has invalid {key}" + ) + for key in ( + "cuvs_java_jar_sha256", + "cuvs_lucene_jar_sha256", + ): + value = build_runtime_artifacts[key] + if ( + not isinstance(value, str) + or re.fullmatch(r"[0-9a-f]{64}", value) is None + ): + raise RuntimeError( + f"Lucene index manifest has invalid {key}" + ) + expected_without_segment_count = { + "schema_version": _MANIFEST_SCHEMA, + "algorithm": self.algorithm, + "codec": self.codec, + "dataset": dict(dataset_identity), + "build_runtime_artifacts": dict(build_runtime_artifacts), + } + actual_without_segment_count = dict(payload) + actual_without_segment_count.pop("segment_count", None) + if actual_without_segment_count != expected_without_segment_count: + raise RuntimeError( + "Existing Lucene index does not match this dataset and configuration; " + "rerun with --force" + ) + + @staticmethod + def _validate_manifest_segment_count( + payload: Mapping[str, Any], verification: Mapping[str, Any] + ) -> None: + stored = payload["segment_count"] + observed = verification.get("segment_count") + if not _is_manifest_integer(observed, minimum=1): + raise RuntimeError( + "Lucene index verification returned an invalid segment count" + ) + if stored != observed: + raise RuntimeError( + "Lucene index manifest segment count does not match the " + f"physical index: {stored} != {observed}; rerun with --force" + ) + + def _search_dataset_identity( + self, dataset: Dataset, payload: Mapping[str, Any] + ) -> tuple[dict[str, Any], int, int]: + stored_identity = payload.get("dataset") + if not isinstance(stored_identity, Mapping): + raise RuntimeError("Lucene index manifest has no dataset identity") + if dataset.base_file: + return _file_backed_dataset_identity( + dataset, + stored_identity, + maximum_dimensions=self.maximum_dimensions, + ) + _normalize_subset_size(dataset) + vectors = _validate_vectors( + dataset.training_vectors, + "training vectors", + maximum_dimensions=self.maximum_dimensions, + ) + return ( + _dataset_identity(dataset, vectors), + int(vectors.shape[0]), + int(vectors.shape[1]), + ) + + def _verification_metadata( + self, + runtime: LuceneRuntime, + path: Path, + vector_count: int, + dimensions: int, + ) -> dict[str, Any]: + verification = runtime.index_verifier.verify( + path, + expected_codec=self.codec, + expected_vector_count=vector_count, + expected_dimensions=dimensions, + ) + metadata = verification.metadata() + metadata["build_route_policy"] = _BUILD_ROUTE_POLICY_BY_ALGORITHM[ + self.algorithm + ] + if self.codec == CAGRA_CODEC: + metadata.update( + runtime.cagra_verifier.verify( + path, + expected_vector_count=vector_count, + expected_dimensions=dimensions, + ).metadata() + ) + return metadata + + @staticmethod + def _install_staged_index(staged: Path, destination: Path) -> str | None: + """Publish a validated index and report non-fatal backup cleanup.""" + if not destination.exists(): + staged.rename(destination) + return None + backup = destination.parent / ( + f".{destination.name}.backup-{uuid.uuid4().hex}" + ) + destination.rename(backup) + try: + staged.rename(destination) + except BaseException as install_error: + try: + backup.rename(destination) + except BaseException as restore_error: + restore_note = ( + "Failed to restore the previous index; its recoverable " + f"backup is {backup}: {type(restore_error).__name__}: " + f"{restore_error}" + ) + if isinstance( + restore_error, (KeyboardInterrupt, SystemExit) + ) and isinstance(install_error, Exception): + restore_error.add_note( + "Index installation first failed: " + f"{type(install_error).__name__}: {install_error}. " + f"The previous index remains in {backup}." + ) + raise + install_error.add_note(restore_note) + raise + try: + shutil.rmtree(backup) + except OSError as cleanup_error: + return ( + f"Published the new index, but could not remove old backup " + f"{backup}: {type(cleanup_error).__name__}: {cleanup_error}" + ) + return None + + def build( + self, + dataset: Dataset, + indexes: list[IndexConfig], + force: bool = False, + dry_run: bool = False, + ) -> BuildResult: + index = self._index(indexes) + try: + _validate_metric(dataset) + path = self._index_path(index) + if dry_run: + return BuildResult( + index_path=str(path), + build_time_seconds=0.0, + index_size_bytes=0, + algorithm=self.algorithm, + build_params=dict(index.build_param), + metadata={ + "dry_run": True, + "codec": self.codec, + "group": self.group, + "index_name": index.name, + }, + ) + if path.exists() and not force: + payload = _read_manifest(path) + dataset_identity, vector_count, dimensions = ( + self._search_dataset_identity(dataset, payload) + ) + self._validate_manifest_identity(payload, dataset_identity) + runtime = self._get_runtime() + runtime.verify_artifacts() + metadata = self._verification_metadata( + runtime, + path, + vector_count, + dimensions, + ) + self._validate_manifest_segment_count(payload, metadata) + return BuildResult( + index_path=str(path), + build_time_seconds=0.0, + index_size_bytes=_index_size(path), + algorithm=self.algorithm, + build_params=dict(index.build_param), + metadata={ + "skipped": True, + "codec": self.codec, + "group": self.group, + "index_name": index.name, + **_artifact_metadata( + "build_runtime", + payload["build_runtime_artifacts"], + ), + **metadata, + }, + ) + backend_build_started = time.perf_counter_ns() + dataset_prepare_started = time.perf_counter_ns() + _normalize_subset_size(dataset) + vectors = _validate_vectors( + dataset.training_vectors, + "training vectors", + maximum_dimensions=self.maximum_dimensions, + ) + dataset_prepare_ns = ( + time.perf_counter_ns() - dataset_prepare_started + ) + if path.exists() and (path.is_symlink() or not path.is_dir()): + raise ValueError( + f"Lucene index path must be a directory: {path}" + ) + path.parent.mkdir(parents=True, exist_ok=True) + runtime_setup_started = time.perf_counter_ns() + runtime = self._get_runtime() + runtime_setup_ns = time.perf_counter_ns() - runtime_setup_started + artifact_validation_started = time.perf_counter_ns() + runtime.verify_artifacts() + artifact_validation_ns = ( + time.perf_counter_ns() - artifact_validation_started + ) + staged = Path( + tempfile.mkdtemp( + prefix=f".{path.name}.build-", dir=path.parent + ) + ) + try: + build_started = time.perf_counter_ns() + runtime_build = runtime.build_index( + staged, vectors, self.codec + ) + build_elapsed_ns = time.perf_counter_ns() - build_started + + validation_started = time.perf_counter_ns() + metadata = self._verification_metadata( + runtime, + staged, + int(vectors.shape[0]), + int(vectors.shape[1]), + ) + observed_segments = int(metadata["segment_count"]) + if runtime_build.segment_count != observed_segments: + raise RuntimeError( + "Lucene build reported a different segment count from " + f"the committed index: {runtime_build.segment_count} != " + f"{observed_segments}" + ) + payload = _manifest_payload( + dataset, + vectors, + self.algorithm, + self.codec, + observed_segments, + runtime.artifact_provenance, + ) + _write_manifest(staged, payload) + validation_elapsed_ns = ( + time.perf_counter_ns() - validation_started + ) + + index_size_started = time.perf_counter_ns() + index_size_bytes = _index_size(staged) + index_size_elapsed_ns = ( + time.perf_counter_ns() - index_size_started + ) + runtime_build_metadata = _runtime_build_timing_metadata( + runtime_build + ) + + install_started = time.perf_counter_ns() + cleanup_warning = self._install_staged_index(staged, path) + install_elapsed_ns = time.perf_counter_ns() - install_started + except BaseException: + shutil.rmtree(staged, ignore_errors=True) + raise + lifecycle_metadata: dict[str, Any] = { + "dataset_load_validate_seconds": _nanoseconds_to_seconds( + dataset_prepare_ns + ), + "runtime_setup_seconds": _nanoseconds_to_seconds( + runtime_setup_ns + ), + "artifact_validation_seconds": _nanoseconds_to_seconds( + artifact_validation_ns + ), + "index_build_call_seconds": _nanoseconds_to_seconds( + build_elapsed_ns + ), + "index_validation_manifest_seconds": ( + _nanoseconds_to_seconds(validation_elapsed_ns) + ), + "index_install_seconds": _nanoseconds_to_seconds( + install_elapsed_ns + ), + "index_size_measurement_seconds": _nanoseconds_to_seconds( + index_size_elapsed_ns + ), + "backend_build_total_seconds": _nanoseconds_to_seconds( + time.perf_counter_ns() - backend_build_started + ), + # Retain the original names for existing result consumers. + "validation_time_seconds": _nanoseconds_to_seconds( + validation_elapsed_ns + ), + "install_time_seconds": _nanoseconds_to_seconds( + install_elapsed_ns + ), + **runtime_build_metadata, + } + if cleanup_warning is not None: + lifecycle_metadata["cleanup_warning"] = cleanup_warning + return BuildResult( + index_path=str(path), + build_time_seconds=_nanoseconds_to_seconds(build_elapsed_ns), + index_size_bytes=index_size_bytes, + algorithm=self.algorithm, + build_params=dict(index.build_param), + metadata={ + "codec": self.codec, + "group": self.group, + "index_name": index.name, + "pylucene_version": runtime.pylucene_version, + **_artifact_metadata( + "build_runtime", runtime.artifact_provenance + ), + **lifecycle_metadata, + **metadata, + }, + ) + except Exception as error: + return self._failure_build(index, error) + + @staticmethod + def _validate_hits( + result: RuntimeSearchResult, k: int, expected_query_count: int + ) -> None: + if len(result.hits) != expected_query_count: + raise RuntimeError( + f"Lucene returned results for {len(result.hits)} queries, " + f"expected {expected_query_count}" + ) + for query_number, hits in enumerate(result.hits): + if len(hits) != k: + raise RuntimeError( + f"Lucene query {query_number} returned {len(hits)} hits, expected {k}" + ) + ids = [hit.document_id for hit in hits] + if len(ids) != len(set(ids)): + raise RuntimeError( + f"Lucene query {query_number} returned duplicate IDs" + ) + if any( + identifier < 0 or identifier >= result.document_count + for identifier in ids + ): + raise RuntimeError( + f"Lucene query {query_number} returned an invalid ID" + ) + + def _successful_search( + self, + runtime_result: RuntimeSearchResult, + parameters: dict[str, int], + *, + k: int, + batch_size: int, + expected_query_count: int, + runtime: LuceneRuntime, + metadata: Mapping[str, Any], + ) -> SearchResult: + self._validate_hits(runtime_result, k, expected_query_count) + timing_metadata, latency_samples_ns = _runtime_search_timing_metadata( + runtime_result, expected_query_count + ) + conversion_started = time.perf_counter_ns() + neighbors = np.asarray( + [ + [hit.document_id for hit in hits] + for hits in runtime_result.hits + ], + dtype=np.int64, + ) + distances = np.asarray( + [ + [_score_to_squared_euclidean(hit.score) for hit in hits] + for hits in runtime_result.hits + ], + dtype=np.float32, + ) + result_array_conversion_ms = _nanoseconds_to_milliseconds( + time.perf_counter_ns() - conversion_started + ) + elapsed_ms = _nanoseconds_to_milliseconds(sum(latency_samples_ns)) + query_count = len(runtime_result.hits) + query_corpus_seconds = _nanoseconds_to_seconds( + runtime_result.timing.query_corpus_wall_ns + ) + if elapsed_ms < 0.0 or query_corpus_seconds <= 0.0: + raise RuntimeError("Lucene returned invalid query timing data") + latency_samples_ms = ( + np.asarray(latency_samples_ns, dtype=np.float64) / 1_000_000.0 + ) + return SearchResult( + neighbors=neighbors, + distances=distances, + search_time_ms=elapsed_ms, + queries_per_second=(query_count / query_corpus_seconds), + recall=0.0, + algorithm=self.algorithm, + search_params=[ + {} if self.algorithm == CAGRA_ALGORITHM else parameters + ], + latency_percentiles={ + "p50": float(np.percentile(latency_samples_ms, 50)), + "p95": float(np.percentile(latency_samples_ms, 95)), + "p99": float(np.percentile(latency_samples_ms, 99)), + }, + metadata={ + "codec": self.codec, + "pylucene_version": runtime.pylucene_version, + **_artifact_metadata( + "search_runtime", runtime.artifact_provenance + ), + "timing_contract_version": _TIMING_CONTRACT_VERSION, + "latency_scope": "client_query", + "throughput_scope": "query_corpus_wall", + "sample_unit": "query", + "execution_model": "serial_single_query", + "latency_seconds": float(latency_samples_ms.mean()) / 1000.0, + "requested_batch_size": batch_size, + "batch_size": batch_size, + "effective_search_batch_size": 1, + "query_count": query_count, + "warmup_query_count": 0, + "result_array_conversion_ms": result_array_conversion_ms, + **timing_metadata, + **metadata, + }, + ) + + def search( + self, + dataset: Dataset, + indexes: list[IndexConfig], + k: int, + batch_size: int = 10000, + mode: str = "latency", + force: bool = False, + search_threads: Optional[int] = None, + dry_run: bool = False, + ) -> list[SearchResult]: + index = self._index(indexes) + try: + _validate_metric(dataset) + if type(k) is not int or k < 1: + raise ValueError("k must be a positive integer") + if ( + isinstance(batch_size, bool) + or not isinstance(batch_size, Integral) + or batch_size < 1 + ): + raise ValueError("batch_size must be a positive integer") + batch_size = int(batch_size) + if mode != "latency": + raise ValueError( + "The initial Lucene backend supports only latency mode" + ) + if search_threads is not None: + raise ValueError( + "The initial Lucene backend does not yet support " + "search_threads" + ) + path = self._index_path(index) + validated_parameters = [ + _search_parameters(self.algorithm, parameters, k) + for parameters in index.search_params + ] + if dry_run: + return [ + SearchResult( + neighbors=np.empty((0, k), dtype=np.int64), + distances=np.empty((0, k), dtype=np.float32), + search_time_ms=0.0, + queries_per_second=0.0, + recall=0.0, + algorithm=self.algorithm, + search_params=[], + metadata={ + "dry_run": True, + "codec": self.codec, + "group": self.group, + "index_name": index.name, + }, + ) + ] + backend_search_started = time.perf_counter_ns() + input_prepare_started = time.perf_counter_ns() + queries = _validate_vectors( + dataset.query_vectors, + "query vectors", + maximum_dimensions=self.maximum_dimensions, + ) + query_input_prepare_ns = ( + time.perf_counter_ns() - input_prepare_started + ) + identity_validation_started = time.perf_counter_ns() + if not path.is_dir(): + raise FileNotFoundError(f"Lucene index does not exist: {path}") + payload = _read_manifest(path) + dataset_identity, vector_count, dimensions = ( + self._search_dataset_identity(dataset, payload) + ) + self._validate_manifest_identity(payload, dataset_identity) + if int(queries.shape[1]) != dimensions: + raise ValueError( + "query vector dimensions do not match the indexed dataset: " + f"{queries.shape[1]} != {dimensions}" + ) + index_identity_validation_ns = ( + time.perf_counter_ns() - identity_validation_started + ) + runtime_setup_started = time.perf_counter_ns() + runtime = self._get_runtime() + runtime_setup_ns = time.perf_counter_ns() - runtime_setup_started + artifact_validation_started = time.perf_counter_ns() + runtime.verify_artifacts() + artifact_validation_ns = ( + time.perf_counter_ns() - artifact_validation_started + ) + index_verification_started = time.perf_counter_ns() + metadata = self._verification_metadata( + runtime, path, vector_count, dimensions + ) + self._validate_manifest_segment_count(payload, metadata) + index_verification_ns = ( + time.perf_counter_ns() - index_verification_started + ) + + results = [] + for parameters in validated_parameters: + search_plan_started = time.perf_counter_ns() + prewarm = _prewarm_index_files(path) + runtime_result = runtime.search_index( + path, + queries, + k=k, + num_candidates=parameters["num_candidates"], + ) + result = self._successful_search( + runtime_result, + parameters, + k=k, + batch_size=batch_size, + expected_query_count=int(queries.shape[0]), + runtime=runtime, + metadata={ + "mode": mode, + "group": self.group, + "index_name": index.name, + **_artifact_metadata( + "build_runtime", + payload["build_runtime_artifacts"], + ), + "expected_search_route": ( + _SEARCH_ROUTE_BY_ALGORITHM[self.algorithm] + ), + "query_input_prepare_ms": ( + _nanoseconds_to_milliseconds( + query_input_prepare_ns + ) + ), + "index_identity_validation_ms": ( + _nanoseconds_to_milliseconds( + index_identity_validation_ns + ) + ), + "runtime_setup_ms": _nanoseconds_to_milliseconds( + runtime_setup_ns + ), + "artifact_validation_ms": ( + _nanoseconds_to_milliseconds( + artifact_validation_ns + ) + ), + "index_verification_ms": ( + _nanoseconds_to_milliseconds(index_verification_ns) + ), + "cache_policy": "sequential_read_all_index_files", + "index_prewarm_ms": _nanoseconds_to_milliseconds( + prewarm.wall_ns + ), + "index_prewarm_bytes": prewarm.bytes_read, + "index_prewarm_file_count": prewarm.file_count, + **metadata, + }, + ) + result.metadata["search_plan_total_ms"] = ( + _nanoseconds_to_milliseconds( + time.perf_counter_ns() - search_plan_started + ) + ) + results.append(result) + backend_search_invocation_total_ms = _nanoseconds_to_milliseconds( + time.perf_counter_ns() - backend_search_started + ) + for result in results: + result.metadata["backend_search_invocation_total_ms"] = ( + backend_search_invocation_total_ms + ) + return results + except Exception as error: + result_width = k if type(k) is int and k > 0 else 0 + return [ + SearchResult( + neighbors=np.empty((0, result_width), dtype=np.int64), + distances=np.empty((0, result_width), dtype=np.float32), + search_time_ms=0.0, + queries_per_second=0.0, + recall=0.0, + algorithm=self.algorithm, + search_params=[], + success=False, + error_message=_format_exception(error), + metadata={ + "group": self.group, + "index_name": index.name, + }, + ) + ] + + +def register() -> None: + """Register the Lucene backend and loader without importing PyLucene.""" + from .registry import ( + _CONFIG_LOADER_REGISTRY, + get_registry, + register_backend, + register_config_loader, + ) + + registry = get_registry() + if not registry.is_registered("lucene"): + register_backend("lucene", LuceneBackend) + if "lucene" not in _CONFIG_LOADER_REGISTRY: + register_config_loader("lucene", LuceneConfigLoader) diff --git a/python/cuvs_bench/cuvs_bench/config/algos/lucene_accelerated_hnsw.yaml b/python/cuvs_bench/cuvs_bench/config/algos/lucene_accelerated_hnsw.yaml new file mode 100644 index 0000000000..fbecef097a --- /dev/null +++ b/python/cuvs_bench/cuvs_bench/config/algos/lucene_accelerated_hnsw.yaml @@ -0,0 +1,13 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +name: lucene_accelerated_hnsw +groups: + base: + build: + codec: ["Lucene101AcceleratedHNSWCodec"] + search: {} + test: + build: + codec: ["Lucene101AcceleratedHNSWCodec"] + search: {} diff --git a/python/cuvs_bench/cuvs_bench/config/algos/lucene_cpu_hnsw.yaml b/python/cuvs_bench/cuvs_bench/config/algos/lucene_cpu_hnsw.yaml new file mode 100644 index 0000000000..34c3922187 --- /dev/null +++ b/python/cuvs_bench/cuvs_bench/config/algos/lucene_cpu_hnsw.yaml @@ -0,0 +1,13 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +name: lucene_cpu_hnsw +groups: + base: + build: + codec: ["Lucene101"] + search: {} + test: + build: + codec: ["Lucene101"] + search: {} diff --git a/python/cuvs_bench/cuvs_bench/config/algos/lucene_cuvs_cagra.yaml b/python/cuvs_bench/cuvs_bench/config/algos/lucene_cuvs_cagra.yaml new file mode 100644 index 0000000000..72f6b162bd --- /dev/null +++ b/python/cuvs_bench/cuvs_bench/config/algos/lucene_cuvs_cagra.yaml @@ -0,0 +1,13 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +name: lucene_cuvs_cagra +groups: + base: + build: + codec: ["CuVS2510GPUSearchCodec"] + search: {} + test: + build: + codec: ["CuVS2510GPUSearchCodec"] + search: {} diff --git a/python/cuvs_bench/cuvs_bench/tests/_lucene_test_support.py b/python/cuvs_bench/cuvs_bench/tests/_lucene_test_support.py new file mode 100644 index 0000000000..60987a3bb6 --- /dev/null +++ b/python/cuvs_bench/cuvs_bench/tests/_lucene_test_support.py @@ -0,0 +1,407 @@ +# +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# + +"""Test doubles and builders shared by the Lucene backend unit tests.""" + +from __future__ import annotations + +from pathlib import Path +from typing import Any, Mapping + +import numpy as np +import pytest + +from cuvs_bench.backends._lucene_runtime import ( + ACCELERATED_HNSW_CODEC, + CAGRA_CODEC, + CPU_HNSW_CODEC, + CagraVerification, + DIRECT_PYLUCENE_DISPATCH, + LuceneIndexVerification, + QueryTiming, + RuntimeBuildResult, + RuntimeBuildTiming, + RuntimeSearchResult, + RuntimeSearchTiming, + SearchHit, + TIMED_BRIDGE_PYLUCENE_DISPATCH, +) +from cuvs_bench.backends._lucene_runtime_config import maven_artifact_version +from cuvs_bench.backends.base import Dataset +from cuvs_bench.backends.lucene import ( + ACCELERATED_HNSW_ALGORITHM, + CAGRA_ALGORITHM, + CPU_HNSW_ALGORITHM, + LuceneBackend, +) +from cuvs_bench.orchestrator.config_loaders import IndexConfig + + +ALGORITHM_CASES = ( + pytest.param( + CPU_HNSW_ALGORITHM, CPU_HNSW_CODEC, "cpu_hnsw", id="cpu-hnsw" + ), + pytest.param( + ACCELERATED_HNSW_ALGORITHM, + ACCELERATED_HNSW_CODEC, + "cpu_hnsw", + id="accelerated-hnsw", + ), + pytest.param(CAGRA_ALGORITHM, CAGRA_CODEC, "gpu_cagra", id="cagra"), +) +_ARTIFACT_VERSION = maven_artifact_version() +_FAKE_ARTIFACT_PROVENANCE = { + "cuvs_java_coordinates": ( + f"com.nvidia.cuvs:cuvs-java:{_ARTIFACT_VERSION}" + ), + "cuvs_java_jar_path": "/artifacts/cuvs-java.jar", + "cuvs_java_jar_sha256": "a" * 64, + "cuvs_lucene_coordinates": ( + f"com.nvidia.cuvs.lucene:cuvs-lucene:{_ARTIFACT_VERSION}" + ), + "cuvs_lucene_jar_path": "/artifacts/cuvs-lucene.jar", + "cuvs_lucene_jar_sha256": "b" * 64, +} + + +class RecordingIndexVerifier: + """Return the physical index facts recorded by the runtime substitute.""" + + def __init__(self, runtime: "RecordingRuntime") -> None: + self.runtime = runtime + self.calls: list[tuple[Path, str, int, int]] = [] + + def verify( + self, + index_path: Path, + *, + expected_codec: str, + expected_vector_count: int, + expected_dimensions: int, + ) -> LuceneIndexVerification: + self.calls.append( + ( + index_path, + expected_codec, + expected_vector_count, + expected_dimensions, + ) + ) + if self.runtime.verification_error is not None: + raise self.runtime.verification_error + return LuceneIndexVerification( + codec=expected_codec, + segment_count=self.runtime.segment_count, + field_count=1, + vector_count=expected_vector_count, + dimensions=expected_dimensions, + ) + + +class RecordingCagraVerifier: + """Return deterministic persisted-path evidence and retain each request.""" + + def __init__(self) -> None: + self.calls: list[tuple[Path, int, int]] = [] + + def verify( + self, + index_path: Path, + *, + expected_vector_count: int, + expected_dimensions: int, + ) -> CagraVerification: + self.calls.append( + (index_path, expected_vector_count, expected_dimensions) + ) + return CagraVerification( + segment_count=1, + field_count=1, + vector_count=expected_vector_count, + dimensions=expected_dimensions, + ) + + +class RecordingRuntime: + """Small in-process substitute for the JVM boundary.""" + + pylucene_version = "10.2.0" + + def __init__(self) -> None: + self.artifact_provenance: dict[str, str] = {} + self.index_verifier = RecordingIndexVerifier(self) + self.cagra_verifier = RecordingCagraVerifier() + self.build_calls: list[tuple[Path, np.ndarray, str]] = [] + self.search_calls: list[dict[str, Any]] = [] + self.build_error: Exception | None = None + self.search_error: Exception | None = None + self.verification_error: Exception | None = None + self.search_result: RuntimeSearchResult | None = None + self.document_count = 0 + self.dimensions = 0 + self.segment_count = 1 + self.artifact_verification_count = 0 + + def verify_artifacts(self) -> None: + self.artifact_verification_count += 1 + + def build_index( + self, index_path: Path, vectors: np.ndarray, codec_name: str + ) -> RuntimeBuildResult: + self.build_calls.append((index_path, vectors.copy(), codec_name)) + if self.build_error is not None: + raise self.build_error + self.document_count, self.dimensions = vectors.shape + (index_path / "segments.fake").write_text(codec_name, encoding="utf-8") + return RuntimeBuildResult( + segment_count=1, + timing=RuntimeBuildTiming( + directory_open_ns=100_000, + writer_setup_ns=200_000, + document_ingest_ns=300_000, + writer_commit_close_ns=400_000, + post_build_reader_ns=500_000, + directory_close_ns=600_000, + runtime_build_wall_ns=2_100_000, + ), + ) + + def search_index( + self, + index_path: Path, + queries: np.ndarray, + *, + k: int, + num_candidates: int, + ) -> RuntimeSearchResult: + self.search_calls.append( + { + "index_path": index_path, + "queries": queries.copy(), + "k": k, + "num_candidates": num_candidates, + } + ) + if self.search_error is not None: + raise self.search_error + if self.search_result is not None: + return self.search_result + hits = [ + [ + SearchHit( + document_id=rank, + score=1.0 / (1.0 + float(rank)), + ) + for rank in range(k) + ] + for _query in queries + ] + query_timings = tuple( + QueryTiming( + query_prepare_ns=100_000, + pylucene_search_dispatch_ns=2_000_000, + java_index_searcher_search_ns=( + 1_500_000 if self.artifact_provenance else None + ), + result_materialization_ns=200_000, + client_query_ns=2_500_000, + ) + for _query in queries + ) + return RuntimeSearchResult( + hits=hits, + timing=RuntimeSearchTiming( + directory_open_ns=250_000, + reader_searcher_setup_ns=750_000, + query_corpus_wall_ns=3_000_000 * len(queries), + reader_close_ns=250_000, + directory_close_ns=250_000, + runtime_plan_wall_ns=3_000_000 * len(queries) + 1_500_000, + search_dispatch_kind=( + TIMED_BRIDGE_PYLUCENE_DISPATCH + if self.artifact_provenance + else DIRECT_PYLUCENE_DISPATCH + ), + queries=query_timings, + ), + document_count=self.document_count, + dimensions=self.dimensions, + ) + + +class RecordingRuntimeFactory: + """Inject one runtime without importing or initializing PyLucene.""" + + def __init__(self, runtime: RecordingRuntime) -> None: + self.runtime = runtime + self.calls: list[dict[str, Any]] = [] + + def __call__(self, config: Mapping[str, Any]) -> RecordingRuntime: + self.calls.append(dict(config)) + return self.runtime + + +def _runtime_search_result( + hits: list[list[SearchHit]], + *, + search_dispatch_ns: tuple[int, ...] | None = None, + java_search_ns: tuple[int, ...] | None = None, + document_count: int = 4, + dimensions: int = 2, + query_corpus_wall_ns: int | None = None, +) -> RuntimeSearchResult: + """Create an explicit per-query timing result for backend tests.""" + query_count = len(hits) + query_prepare_ns = 100_000 + result_materialization_ns = 200_000 + unclassified_client_overhead_ns = 200_000 + directory_open_ns = 250_000 + reader_searcher_setup_ns = 750_000 + reader_close_ns = 250_000 + directory_close_ns = 250_000 + dispatch_samples = ( + search_dispatch_ns + if search_dispatch_ns is not None + else (2_000_000,) * query_count + ) + if len(dispatch_samples) != query_count: + raise ValueError("search_dispatch_ns must contain one value per query") + if java_search_ns is not None and len(java_search_ns) != query_count: + raise ValueError("java_search_ns must contain one value per query") + query_timings = tuple( + QueryTiming( + query_prepare_ns=query_prepare_ns, + pylucene_search_dispatch_ns=dispatch_samples[index], + java_index_searcher_search_ns=( + java_search_ns[index] if java_search_ns is not None else None + ), + result_materialization_ns=result_materialization_ns, + client_query_ns=( + query_prepare_ns + + dispatch_samples[index] + + result_materialization_ns + + unclassified_client_overhead_ns + ), + ) + for index in range(query_count) + ) + corpus_ns = ( + query_corpus_wall_ns + if query_corpus_wall_ns is not None + else sum(item.client_query_ns for item in query_timings) + ) + return RuntimeSearchResult( + hits=hits, + timing=RuntimeSearchTiming( + directory_open_ns=directory_open_ns, + reader_searcher_setup_ns=reader_searcher_setup_ns, + query_corpus_wall_ns=corpus_ns, + reader_close_ns=reader_close_ns, + directory_close_ns=directory_close_ns, + runtime_plan_wall_ns=( + directory_open_ns + + reader_searcher_setup_ns + + corpus_ns + + reader_close_ns + + directory_close_ns + ), + search_dispatch_kind=( + TIMED_BRIDGE_PYLUCENE_DISPATCH + if java_search_ns is not None + else DIRECT_PYLUCENE_DISPATCH + ), + queries=query_timings, + ), + document_count=document_count, + dimensions=dimensions, + ) + + +def _dataset(*, offset: float = 0.0) -> Dataset: + training_vectors = np.asarray( + [ + [0.0 + offset, 0.0], + [1.0, 0.0], + [0.0, 2.0], + [3.0, 4.0], + ], + dtype=np.float32, + ) + return Dataset( + name="tiny-l2", + training_vectors=training_vectors, + query_vectors=training_vectors[:2].copy(), + distance_metric="euclidean", + ) + + +def _dataset_with_dimensions(dimensions: int) -> Dataset: + vectors = np.zeros((2, dimensions), dtype=np.float32) + return Dataset( + name=f"dimensions-{dimensions}", + training_vectors=vectors, + query_vectors=vectors[:1].copy(), + distance_metric="euclidean", + ) + + +def _write_fbin(path: Path, vectors: np.ndarray) -> None: + rows, dimensions = vectors.shape + path.write_bytes( + np.asarray([rows, dimensions], dtype=np.uint32).tobytes() + + np.ascontiguousarray(vectors, dtype=np.float32).tobytes() + ) + + +def _artifact_pair(tmp_path: Path) -> tuple[Path, Path]: + java_jar = tmp_path / "cuvs-java.jar" + lucene_jar = tmp_path / "cuvs-lucene.jar" + java_jar.touch() + lucene_jar.touch() + return java_jar, lucene_jar + + +def _backend_and_index( + tmp_path: Path, + algorithm: str, + runtime: RecordingRuntime, + *, + search_params: list[dict[str, Any]] | None = None, +) -> tuple[LuceneBackend, IndexConfig, RecordingRuntimeFactory]: + codec = { + CPU_HNSW_ALGORITHM: CPU_HNSW_CODEC, + ACCELERATED_HNSW_ALGORITHM: ACCELERATED_HNSW_CODEC, + CAGRA_ALGORITHM: CAGRA_CODEC, + }[algorithm] + index_root = tmp_path / "indexes" + config: dict[str, Any] = { + "name": algorithm, + "algo": algorithm, + "codec": codec, + "group": "test", + "index_root": str(index_root), + "requires_cuvs": algorithm != CPU_HNSW_ALGORITHM, + } + if algorithm != CPU_HNSW_ALGORITHM: + runtime.artifact_provenance = dict(_FAKE_ARTIFACT_PROVENANCE) + java_jar, lucene_jar = _artifact_pair(tmp_path) + (tmp_path / "libcuvs_c.so").touch() + config.update( + { + "cuvs_java_jar": str(java_jar), + "cuvs_lucene_jar": str(lucene_jar), + "java_library_path": str(tmp_path), + } + ) + factory = RecordingRuntimeFactory(runtime) + backend = LuceneBackend(config, runtime_factory=factory) + index = IndexConfig( + name=algorithm, + algo=algorithm, + build_param={"codec": codec}, + search_params=search_params if search_params is not None else [{}], + file=str(index_root / algorithm), + ) + return backend, index, factory diff --git a/python/cuvs_bench/cuvs_bench/tests/test_lucene_backend_contract.py b/python/cuvs_bench/cuvs_bench/tests/test_lucene_backend_contract.py new file mode 100644 index 0000000000..fcbb62f78f --- /dev/null +++ b/python/cuvs_bench/cuvs_bench/tests/test_lucene_backend_contract.py @@ -0,0 +1,72 @@ +# +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Registration and packaged-configuration contracts for the Lucene backend.""" + +from importlib import resources + +import yaml + +from cuvs_bench.backends.base import BuildResult +from cuvs_bench.backends.lucene import ( + CAGRA_ALGORITHM, + LuceneBackend, + LuceneConfigLoader, + register, +) +from cuvs_bench.backends.registry import get_backend_class, get_config_loader + + +def test_registration_exposes_the_backend_and_config_loader() -> None: + register() + + backend_class = get_backend_class("lucene") + assert backend_class is LuceneBackend + assert backend_class.default_algorithm == CAGRA_ALGORITHM + assert get_config_loader("lucene") is LuceneConfigLoader + + +def test_algorithm_configs_define_each_supported_codec() -> None: + expected_codecs = { + "lucene_accelerated_hnsw": "Lucene101AcceleratedHNSWCodec", + "lucene_cpu_hnsw": "Lucene101", + "lucene_cuvs_cagra": "CuVS2510GPUSearchCodec", + } + algorithm_resources = resources.files("cuvs_bench.config.algos") + + for algorithm, codec in expected_codecs.items(): + config = yaml.safe_load( + algorithm_resources.joinpath(f"{algorithm}.yaml").read_text() + ) + assert config["name"] == algorithm + assert set(config["groups"]) == {"base", "test"} + for group in config["groups"].values(): + assert group == { + "build": {"codec": [codec]}, + "search": {}, + } + + +def test_failed_results_are_fatal_for_the_lucene_backend() -> None: + successful = BuildResult( + index_path="", + build_time_seconds=0.0, + index_size_bytes=0, + algorithm=CAGRA_ALGORITHM, + build_params={}, + ) + failed = BuildResult( + index_path="", + build_time_seconds=0.0, + index_size_bytes=0, + algorithm=CAGRA_ALGORITHM, + build_params={}, + success=False, + error_message="GPU path was unavailable", + ) + + assert LuceneBackend.result_failure_message([successful]) is None + assert LuceneBackend.result_failure_message([successful, failed]) == ( + "GPU path was unavailable" + ) diff --git a/python/cuvs_bench/cuvs_bench/tests/test_lucene_build_search.py b/python/cuvs_bench/cuvs_bench/tests/test_lucene_build_search.py new file mode 100644 index 0000000000..27a9b48aaa --- /dev/null +++ b/python/cuvs_bench/cuvs_bench/tests/test_lucene_build_search.py @@ -0,0 +1,837 @@ +# +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# + +"""Build and search behavior for the opt-in Lucene benchmark backend.""" + +from __future__ import annotations + +import json +from dataclasses import replace +from pathlib import Path + +import numpy as np +import pytest + +from _lucene_test_support import ( + ALGORITHM_CASES, + _FAKE_ARTIFACT_PROVENANCE, + RecordingRuntime, + RecordingRuntimeFactory, + _backend_and_index, + _dataset, + _dataset_with_dimensions, + _runtime_search_result, + _write_fbin, +) +from cuvs_bench.backends._lucene_runtime import ( + RuntimeBuildResult, + RuntimeBuildTiming, + RuntimeSearchResult, + SearchHit, +) +from cuvs_bench.backends.base import Dataset +from cuvs_bench.backends.lucene import ( + ACCELERATED_HNSW_ALGORITHM, + CAGRA_ALGORITHM, + CPU_HNSW_ALGORITHM, + MAX_CAGRA_DIMENSIONS, + MAX_CPU_HNSW_DIMENSIONS, + LuceneBackend, + _file_backed_dataset_identity, + _runtime_build_timing_metadata, + _runtime_search_timing_metadata, +) +from cuvs_bench.orchestrator.config_loaders import IndexConfig + + +@pytest.mark.parametrize(("algorithm", "codec", "path"), ALGORITHM_CASES) +def test_successful_build_reports_the_verified_persisted_index_kind( + tmp_path: Path, algorithm: str, codec: str, path: str +) -> None: + runtime = RecordingRuntime() + backend, index, factory = _backend_and_index(tmp_path, algorithm, runtime) + + result = backend.build(_dataset(), [index]) + + assert result.success, result.error_message + assert result.algorithm == algorithm + assert result.build_params == {"codec": codec} + assert result.metadata["codec"] == codec + expected_persisted_kind = { + CPU_HNSW_ALGORITHM: "cpu_hnsw", + ACCELERATED_HNSW_ALGORITHM: "hnsw", + CAGRA_ALGORITHM: "gpu_cagra_only", + }[algorithm] + assert result.metadata["persisted_index_kind"] == expected_persisted_kind + assert ( + result.metadata["build_route_policy"] + == { + CPU_HNSW_ALGORITHM: "cpu_hnsw", + ACCELERATED_HNSW_ALGORITHM: "gpu_cagra_or_cpu_hnsw_fallback", + CAGRA_ALGORITHM: "gpu_cagra", + }[algorithm] + ) + assert result.metadata["segment_count"] == 1 + assert result.metadata["pylucene_version"] == "10.2.0" + assert result.index_size_bytes > 0 + assert len(factory.calls) == 1 + assert [call[2] for call in runtime.build_calls] == [codec] + if algorithm == CAGRA_ALGORITHM: + [(verified_path, vector_count, dimensions)] = ( + runtime.cagra_verifier.calls + ) + destination = Path(index.file).resolve() + assert verified_path.parent == destination.parent + assert verified_path.name.startswith(f".{destination.name}.build-") + assert not verified_path.exists() + assert (vector_count, dimensions) == (4, 2) + else: + assert runtime.cagra_verifier.calls == [] + + +@pytest.mark.parametrize( + ("algorithm", "dimensions", "maximum_dimensions", "expected_success"), + ( + pytest.param( + CPU_HNSW_ALGORITHM, + 1024, + MAX_CPU_HNSW_DIMENSIONS, + True, + id="cpu-maximum-dimensions", + ), + pytest.param( + CPU_HNSW_ALGORITHM, + 1025, + MAX_CPU_HNSW_DIMENSIONS, + False, + id="cpu-above-maximum-dimensions", + ), + pytest.param( + ACCELERATED_HNSW_ALGORITHM, + 4096, + MAX_CAGRA_DIMENSIONS, + True, + id="accelerated-hnsw-maximum-dimensions", + ), + pytest.param( + ACCELERATED_HNSW_ALGORITHM, + 4097, + MAX_CAGRA_DIMENSIONS, + False, + id="accelerated-hnsw-above-maximum-dimensions", + ), + pytest.param( + CAGRA_ALGORITHM, + 1025, + MAX_CAGRA_DIMENSIONS, + True, + id="cagra-above-cpu-maximum", + ), + pytest.param( + CAGRA_ALGORITHM, + 4096, + MAX_CAGRA_DIMENSIONS, + True, + id="cagra-maximum-dimensions", + ), + pytest.param( + CAGRA_ALGORITHM, + 4097, + MAX_CAGRA_DIMENSIONS, + False, + id="cagra-above-maximum-dimensions", + ), + ), +) +def test_build_enforces_each_algorithms_dimension_limit( + tmp_path: Path, + algorithm: str, + dimensions: int, + maximum_dimensions: int, + expected_success: bool, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index(tmp_path, algorithm, runtime) + + result = backend.build(_dataset_with_dimensions(dimensions), [index]) + + assert result.success is expected_success + if expected_success: + assert runtime.dimensions == dimensions + else: + assert result.error_message == ( + "ValueError: training vectors dimensions must not exceed " + f"{maximum_dimensions}" + ) + assert runtime.build_calls == [] + + +def test_file_backed_identity_rejects_cpu_dimensions_before_hashing( + tmp_path: Path, +) -> None: + vectors = np.zeros((2, MAX_CPU_HNSW_DIMENSIONS + 1), dtype=np.float32) + source = tmp_path / "base.fbin" + _write_fbin(source, vectors) + dataset = _dataset_with_dimensions(2) + dataset.base_file = str(source) + + with pytest.raises(ValueError, match="must not exceed 1024"): + _file_backed_dataset_identity( + dataset, + {}, + maximum_dimensions=MAX_CPU_HNSW_DIMENSIONS, + ) + + +def test_numpy_subset_size_is_normalized_before_manifest_serialization( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + vectors = _dataset().training_vectors + base_file = tmp_path / "base.fbin" + _write_fbin(base_file, vectors) + dataset = Dataset( + name="tiny-l2-subset", + query_vectors=vectors[:2].copy(), + distance_metric="euclidean", + base_file=str(base_file), + metadata={"subset_size": np.int64(3)}, + ) + + first_result = backend.build(dataset, [index]) + second_result = backend.build(dataset, [index]) + + assert first_result.success, first_result.error_message + assert second_result.success, second_result.error_message + assert second_result.metadata["skipped"] is True + assert runtime.document_count == 3 + assert len(runtime.build_calls) == 1 + assert type(dataset.metadata["subset_size"]) is int + manifest = json.loads( + (Path(index.file) / ".cuvs-bench-lucene.json").read_text( + encoding="utf-8" + ) + ) + assert manifest["dataset"]["subset_size"] == 3 + + +def test_boolean_subset_size_is_rejected_before_runtime_initialization( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = _dataset() + dataset.metadata = {"subset_size": True} + + result = backend.build(dataset, [index]) + + assert not result.success + assert result.error_message == ( + "ValueError: subset_size must be a positive integer, got True" + ) + assert factory.calls == [] + assert runtime.build_calls == [] + assert not Path(index.file).exists() + + +@pytest.mark.parametrize(("algorithm", "codec", "path"), ALGORITHM_CASES) +def test_successful_search_reports_path_parameters_and_distances( + tmp_path: Path, algorithm: str, codec: str, path: str +) -> None: + runtime = RecordingRuntime() + search_params = ( + [{"num_candidates": 4}] if algorithm != CAGRA_ALGORITHM else [{}] + ) + backend, index, _factory = _backend_and_index( + tmp_path, + algorithm, + runtime, + search_params=search_params, + ) + build_result = backend.build(_dataset(), [index]) + assert build_result.success, build_result.error_message + + result = backend.search(_dataset(), [index], k=2, batch_size=1)[0] + + assert result.success, result.error_message + np.testing.assert_array_equal(result.neighbors, [[0, 1], [0, 1]]) + np.testing.assert_allclose(result.distances, [[0.0, 1.0], [0.0, 1.0]]) + assert result.metadata["codec"] == codec + assert result.metadata["pylucene_version"] == "10.2.0" + assert result.metadata["latency_seconds"] == 0.0025 + assert result.metadata["timing_contract_version"] == 1 + assert result.metadata["latency_scope"] == "client_query" + assert result.metadata["throughput_scope"] == "query_corpus_wall" + assert result.metadata["sample_unit"] == "query" + assert result.metadata["execution_model"] == "serial_single_query" + assert result.metadata["requested_batch_size"] == 1 + assert result.metadata["effective_search_batch_size"] == 1 + assert result.metadata["query_count"] == 2 + assert result.metadata["pylucene_search_dispatch_count"] == 2 + assert result.metadata["search_dispatch_kind"] == ( + "thin_jar_timing_bridge" + if algorithm != CPU_HNSW_ALGORITHM + else "direct_pylucene" + ) + assert result.metadata["batch_size"] == 1 + assert result.metadata["mode"] == "latency" + assert result.metadata["group"] == "test" + assert result.metadata["index_name"] == algorithm + assert result.metadata["expected_search_route"] == path + assert ( + result.metadata["persisted_index_kind"] + == { + CPU_HNSW_ALGORITHM: "cpu_hnsw", + ACCELERATED_HNSW_ALGORITHM: "hnsw", + CAGRA_ALGORITHM: "gpu_cagra_only", + }[algorithm] + ) + assert result.metadata["segment_count"] == 1 + assert result.metadata["field_count"] == 1 + assert result.metadata["vector_count"] == 4 + assert result.metadata["dimensions"] == 2 + assert result.metadata["index_prewarm_bytes"] > 0 + assert result.metadata["index_prewarm_file_count"] >= 2 + assert result.metadata["java_timing_available"] is ( + algorithm != CPU_HNSW_ALGORITHM + ) + if algorithm != CPU_HNSW_ALGORITHM: + for role in ("build_runtime", "search_runtime"): + for key, value in _FAKE_ARTIFACT_PROVENANCE.items(): + assert result.metadata[f"{role}_{key}"] == value + assert result.search_params == ( + [{"num_candidates": 4}] if algorithm != CAGRA_ALGORITHM else [{}] + ) + assert runtime.search_calls[0]["num_candidates"] == ( + 4 if algorithm != CAGRA_ALGORITHM else 2 + ) + + +def test_search_parameter_sweep_verifies_artifacts_once( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + search_params = [{"num_candidates": 2}, {"num_candidates": 4}] + backend, index, _factory = _backend_and_index( + tmp_path, + CPU_HNSW_ALGORITHM, + runtime, + search_params=search_params, + ) + assert backend.build(_dataset(), [index]).success + verifications_after_build = runtime.artifact_verification_count + + results = backend.search(_dataset(), [index], k=2, batch_size=2) + + assert all(result.success for result in results) + assert len(results) == len(search_params) + assert runtime.artifact_verification_count == verifications_after_build + 1 + assert [ + search_call["num_candidates"] for search_call in runtime.search_calls + ] == [ + 2, + 4, + ] + invocation_totals = { + result.metadata["backend_search_invocation_total_ms"] + for result in results + } + assert len(invocation_totals) == 1 + assert next(iter(invocation_totals)) >= max( + result.metadata["search_plan_total_ms"] for result in results + ) + assert all( + "backend_search_total_ms" not in result.metadata for result in results + ) + + +def test_search_timing_uses_one_sample_per_serial_query( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = _dataset() + assert backend.build(dataset, [index]).success + runtime.search_result = _runtime_search_result( + [ + [SearchHit(0, 1.0), SearchHit(1, 0.5)], + [SearchHit(0, 1.0), SearchHit(1, 0.5)], + ], + search_dispatch_ns=(2_000_000, 8_000_000), + query_corpus_wall_ns=12_000_000, + ) + + result = backend.search(dataset, [index], k=2, batch_size=2)[0] + + assert result.success, result.error_message + assert result.latency_percentiles == { + "p50": 5.5, + "p95": pytest.approx(8.2), + "p99": pytest.approx(8.44), + } + assert result.search_time_ms == 11.0 + assert result.queries_per_second == pytest.approx(2.0 / 0.012) + assert result.metadata["latency_seconds"] == 0.0055 + assert result.metadata["pylucene_search_dispatch_count"] == 2 + assert result.metadata["first_query_pylucene_search_dispatch_ms"] == 2.0 + assert result.metadata["first_query_client_query_ms"] == 2.5 + assert result.metadata["subsequent_query_count"] == 1 + assert ( + result.metadata["subsequent_pylucene_search_dispatch_mean_ms"] == 8.0 + ) + assert result.metadata["subsequent_client_query_mean_ms"] == 8.5 + assert result.metadata["requested_batch_size"] == 2 + assert result.metadata["effective_search_batch_size"] == 1 + + +def test_exact_java_timing_reports_first_and_subsequent_queries_only_when_available() -> ( + None +): + hits = [ + [SearchHit(0, 1.0), SearchHit(1, 0.5)], + [SearchHit(0, 1.0), SearchHit(1, 0.5)], + ] + bridged = _runtime_search_result( + hits, + java_search_ns=(1_000_000, 4_000_000), + ) + direct = _runtime_search_result(hits) + + bridged_metadata, _ = _runtime_search_timing_metadata(bridged, 2) + direct_metadata, _ = _runtime_search_timing_metadata(direct, 2) + + assert bridged_metadata["first_query_java_index_searcher_search_ms"] == 1.0 + assert ( + bridged_metadata["subsequent_java_index_searcher_search_mean_ms"] + == 4.0 + ) + assert "first_query_java_index_searcher_search_ms" not in direct_metadata + assert ( + "subsequent_java_index_searcher_search_mean_ms" not in direct_metadata + ) + + +def test_requested_batch_size_never_groups_lucene_query_samples( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = _dataset() + dataset.query_vectors = dataset.training_vectors[:3].copy() + assert backend.build(dataset, [index]).success + + result = backend.search(dataset, [index], k=2, batch_size=2)[0] + + assert result.success, result.error_message + assert result.metadata["query_count"] == 3 + assert result.metadata["client_query_count"] == 3 + assert result.metadata["pylucene_search_dispatch_count"] == 3 + assert result.metadata["requested_batch_size"] == 2 + assert result.metadata["effective_search_batch_size"] == 1 + assert "batch_count" not in result.metadata + + +def test_search_rejects_missing_per_query_timing_data(tmp_path: Path) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = _dataset() + assert backend.build(dataset, [index]).success + single_query_result = _runtime_search_result( + [[SearchHit(0, 1.0), SearchHit(1, 0.5)]] + ) + runtime.search_result = RuntimeSearchResult( + hits=[ + [SearchHit(0, 1.0), SearchHit(1, 0.5)], + [SearchHit(0, 1.0), SearchHit(1, 0.5)], + ], + timing=single_query_result.timing, + document_count=4, + dimensions=2, + ) + + result = backend.search(dataset, [index], k=2)[0] + + assert not result.success + assert result.error_message == ( + "RuntimeError: Lucene returned timing data for 1 queries, expected 2" + ) + + +def test_build_timing_rejects_phases_longer_than_the_enclosing_wall() -> None: + timing = RuntimeBuildTiming( + directory_open_ns=1, + writer_setup_ns=2, + document_ingest_ns=3, + writer_commit_close_ns=4, + post_build_reader_ns=5, + directory_close_ns=6, + runtime_build_wall_ns=20, + ) + + with pytest.raises(RuntimeError, match="impossible runtime_build_wall"): + _runtime_build_timing_metadata( + RuntimeBuildResult(segment_count=1, timing=timing) + ) + + +def test_search_timing_rejects_query_phases_outside_client_wall() -> None: + result = _runtime_search_result([[SearchHit(0, 1.0), SearchHit(1, 0.5)]]) + query = replace(result.timing.queries[0], client_query_ns=2_000_000) + contradictory = replace( + result, + timing=replace(result.timing, queries=(query,)), + ) + + with pytest.raises(RuntimeError, match=r"impossible client_query\[0\]"): + _runtime_search_timing_metadata(contradictory, 1) + + +def test_search_timing_rejects_client_total_longer_than_corpus_wall() -> None: + result = _runtime_search_result([[SearchHit(0, 1.0), SearchHit(1, 0.5)]]) + contradictory = replace( + result, + timing=replace(result.timing, query_corpus_wall_ns=2_499_999), + ) + + with pytest.raises(RuntimeError, match="impossible query_corpus_wall"): + _runtime_search_timing_metadata(contradictory, 1) + + +def test_search_timing_rejects_plan_phases_longer_than_plan_wall() -> None: + result = _runtime_search_result([[SearchHit(0, 1.0), SearchHit(1, 0.5)]]) + contradictory = replace( + result, + timing=replace(result.timing, runtime_plan_wall_ns=3_999_999), + ) + + with pytest.raises(RuntimeError, match="impossible runtime_plan_wall"): + _runtime_search_timing_metadata(contradictory, 1) + + +@pytest.mark.parametrize( + ("hits", "message"), + ( + pytest.param( + [[SearchHit(0, 1.0), SearchHit(0, 0.5)]] * 2, + "duplicate IDs", + id="duplicate", + ), + pytest.param( + [[SearchHit(0, 1.0), SearchHit(4, 0.5)]] * 2, + "an invalid ID", + id="out-of-range", + ), + ), +) +def test_search_rejects_invalid_hit_identifiers( + tmp_path: Path, hits: list[list[SearchHit]], message: str +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + assert backend.build(_dataset(), [index]).success + runtime.search_result = _runtime_search_result(hits) + + result = backend.search(_dataset(), [index], k=2)[0] + + assert not result.success + assert ( + f"RuntimeError: Lucene query 0 returned {message}" + in result.error_message + ) + + +def test_search_rejects_a_runtime_result_that_drops_a_query( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = _dataset() + assert backend.build(dataset, [index]).success + runtime.search_result = _runtime_search_result( + [[SearchHit(0, 1.0), SearchHit(1, 0.5)]] + ) + + result = backend.search(dataset, [index], k=2)[0] + + assert not result.success + assert result.error_message == ( + "RuntimeError: Lucene returned results for 1 queries, expected 2" + ) + + +def test_cagra_search_above_1024_fails_before_issuing_a_query( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CAGRA_ALGORITHM, runtime + ) + assert backend.build(_dataset(), [index]).success + + result = backend.search(_dataset(), [index], k=1025)[0] + + assert not result.success + assert "CAGRA search supports k <= 1024" in result.error_message + assert runtime.search_calls == [] + + +def test_invalid_top_k_returns_the_validation_error_without_a_secondary_failure( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + + result = backend.search(_dataset(), [index], k=-1)[0] + + assert not result.success + assert result.error_message == "ValueError: k must be a positive integer" + assert result.neighbors.shape == (0, 0) + assert result.distances.shape == (0, 0) + + +@pytest.mark.parametrize("batch_size", (0, 2.5, True, np.bool_(True))) +def test_invalid_batch_size_fails_before_issuing_a_query( + tmp_path: Path, batch_size: object +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + + result = backend.search(_dataset(), [index], k=2, batch_size=batch_size)[0] + + assert not result.success + assert result.error_message == ( + "ValueError: batch_size must be a positive integer" + ) + assert runtime.search_calls == [] + + +def test_numpy_integer_batch_size_is_reported_without_changing_execution( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + assert backend.build(_dataset(), [index]).success + + result = backend.search(_dataset(), [index], k=2, batch_size=np.int64(2))[ + 0 + ] + + assert result.success, result.error_message + assert "batch_size" not in runtime.search_calls[0] + assert result.metadata["requested_batch_size"] == 2 + assert result.metadata["effective_search_batch_size"] == 1 + + +def test_build_failure_is_actionable_and_removes_partial_index( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + runtime.build_error = RuntimeError("GPU device is unavailable") + backend, index, _factory = _backend_and_index( + tmp_path, CAGRA_ALGORITHM, runtime + ) + + result = backend.build(_dataset(), [index]) + + assert not result.success + assert result.error_message == "RuntimeError: GPU device is unavailable" + assert not Path(index.file).exists() + + +def test_build_timing_separates_build_validation_and_index_publication( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + clock = iter( + int(seconds * 1_000_000_000) + for seconds in ( + 0, + 10, + 11, + 12, + 13, + 14, + 15, + 20, + 22, + 30, + 35, + 40, + 41, + 43, + 50, + 60, + ) + ) + + def read_clock() -> int: + return next(clock) + + monkeypatch.setattr( + "cuvs_bench.backends.lucene.time.perf_counter_ns", read_clock + ) + + result = backend.build(_dataset(), [index]) + + assert result.success, result.error_message + assert result.build_time_seconds == 2.0 + assert result.metadata["validation_time_seconds"] == 5.0 + assert result.metadata["install_time_seconds"] == 7.0 + assert result.metadata["dataset_load_validate_seconds"] == 1.0 + assert result.metadata["runtime_setup_seconds"] == 1.0 + assert result.metadata["artifact_validation_seconds"] == 1.0 + assert result.metadata["index_build_call_seconds"] == 2.0 + assert result.metadata["index_validation_manifest_seconds"] == 5.0 + assert result.metadata["index_install_seconds"] == 7.0 + assert result.metadata["index_size_measurement_seconds"] == 1.0 + assert result.metadata["backend_build_total_seconds"] == 60.0 + assert result.metadata["runtime_document_ingest_seconds"] == 0.0003 + assert runtime.artifact_verification_count == 1 + + +def test_search_failure_is_returned_with_context( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + assert backend.build(_dataset(), [index]).success + runtime.search_error = RuntimeError( + "Lucene reader could not open the index" + ) + + result = backend.search(_dataset(), [index], k=2)[0] + + assert not result.success + assert result.error_message == ( + "RuntimeError: Lucene reader could not open the index" + ) + assert result.neighbors.shape == (0, 2) + assert result.distances.shape == (0, 2) + + +def test_failure_diagnostics_include_exception_notes(tmp_path: Path) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + error = RuntimeError("writer failed") + error.add_note("Rollback also failed: disk unavailable") + runtime.build_error = error + + result = backend.build(_dataset(), [index]) + + assert not result.success + assert result.error_message == ( + "RuntimeError: writer failed\nRollback also failed: disk unavailable" + ) + + +def test_throughput_mode_fails_instead_of_reporting_serial_search_as_throughput( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + + result = backend.search(_dataset(), [index], k=2, mode="throughput")[0] + + assert not result.success + assert "supports only latency mode" in result.error_message + assert runtime.search_calls == [] + + +@pytest.mark.parametrize(("algorithm", "codec", "_path"), ALGORITHM_CASES) +def test_dry_runs_do_not_resolve_or_start_the_runtime( + tmp_path: Path, algorithm: str, codec: str, _path: str +) -> None: + runtime = RecordingRuntime() + factory = RecordingRuntimeFactory(runtime) + index_root = tmp_path / "indexes" + backend = LuceneBackend( + { + "name": algorithm, + "algo": algorithm, + "codec": codec, + "group": "test", + "index_root": str(index_root), + "requires_cuvs": algorithm != CPU_HNSW_ALGORITHM, + }, + runtime_factory=factory, + ) + index = IndexConfig( + name=algorithm, + algo=algorithm, + build_param={"codec": codec}, + search_params=[{}], + file=str(index_root / algorithm), + ) + + build_result = backend.build(_dataset(), [index], dry_run=True) + search_result = backend.search(_dataset(), [index], k=2, dry_run=True)[0] + + assert build_result.success + assert build_result.metadata == { + "dry_run": True, + "codec": codec, + "group": "test", + "index_name": algorithm, + } + assert search_result.success + assert search_result.metadata == { + "dry_run": True, + "codec": codec, + "group": "test", + "index_name": algorithm, + } + assert factory.calls == [] + assert not index_root.exists() + + +def test_dry_runs_do_not_materialize_lazy_vectors(tmp_path: Path) -> None: + class UnreadableVectors(Dataset): + @property + def training_vectors(self) -> np.ndarray: + raise AssertionError("dry-run loaded training vectors") + + @property + def query_vectors(self) -> np.ndarray: + raise AssertionError("dry-run loaded query vectors") + + runtime = RecordingRuntime() + backend, index, factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = UnreadableVectors(name="tiny-l2", distance_metric="euclidean") + + assert backend.build(dataset, [index], dry_run=True).success + assert backend.search(dataset, [index], k=2, dry_run=True)[0].success + assert factory.calls == [] diff --git a/python/cuvs_bench/cuvs_bench/tests/test_lucene_lifecycle.py b/python/cuvs_bench/cuvs_bench/tests/test_lucene_lifecycle.py new file mode 100644 index 0000000000..ce8313b2b5 --- /dev/null +++ b/python/cuvs_bench/cuvs_bench/tests/test_lucene_lifecycle.py @@ -0,0 +1,708 @@ +# +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# + +"""Index publication, reuse, and provenance tests for the Lucene backend.""" + +from __future__ import annotations + +import json +from pathlib import Path + +import numpy as np +import pytest + +from _lucene_test_support import ( + RecordingRuntime, + _backend_and_index, + _dataset, + _write_fbin, +) +from cuvs_bench.backends._lucene_runtime import CPU_HNSW_CODEC +from cuvs_bench.backends.base import Dataset +from cuvs_bench.backends.lucene import ( + CAGRA_ALGORITHM, + CPU_HNSW_ALGORITHM, + LuceneBackend, + _prewarm_index_files, + _source_identity, +) +from cuvs_bench.orchestrator.config_loaders import IndexConfig + + +def test_index_prewarm_reads_every_regular_file(tmp_path: Path) -> None: + (tmp_path / "segments_1").write_bytes(b"segments") + (tmp_path / "vectors.vec").write_bytes(b"vector-data") + (tmp_path / "ignored-directory").mkdir() + + timing = _prewarm_index_files(tmp_path) + + assert timing.file_count == 2 + assert timing.bytes_read == len(b"segments") + len(b"vector-data") + assert timing.wall_ns >= 0 + + +def test_index_prewarm_reports_the_file_that_could_not_be_read( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + unreadable = tmp_path / "vectors.vec" + unreadable.write_bytes(b"vector-data") + original_open = Path.open + + def fail_open(path: Path, *args, **kwargs): + if path == unreadable: + raise PermissionError("read denied") + return original_open(path, *args, **kwargs) + + monkeypatch.setattr(Path, "open", fail_open) + + with pytest.raises( + RuntimeError, + match=f"Failed to prewarm Lucene index file: {unreadable}", + ) as failure: + _prewarm_index_files(tmp_path) + + assert isinstance(failure.value.__cause__, PermissionError) + + +def test_index_prewarm_refuses_symbolic_links(tmp_path: Path) -> None: + target = tmp_path / "target" + target.write_bytes(b"index-data") + link = tmp_path / "linked-index-file" + link.symlink_to(target) + + with pytest.raises(RuntimeError, match="refuses symbolic links"): + _prewarm_index_files(tmp_path) + + +def test_force_rebuild_never_removes_a_path_outside_the_configured_root( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, _index, factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + outside = tmp_path / "outside" / "index" + outside.mkdir(parents=True) + sentinel = outside / "keep-me" + sentinel.write_text("preserve", encoding="utf-8") + index = IndexConfig( + name="outside", + algo=CPU_HNSW_ALGORITHM, + build_param={"codec": CPU_HNSW_CODEC}, + search_params=[{}], + file=str(outside), + ) + + result = backend.build(_dataset(), [index], force=True) + + assert not result.success + assert "outside its configured root" in result.error_message + assert sentinel.read_text(encoding="utf-8") == "preserve" + assert factory.calls == [] + + +def test_force_rebuild_rejects_a_symlink_to_a_sibling_index( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + root = Path(backend.config["index_root"]) + root.mkdir(parents=True) + sibling = root / "existing-index" + sibling.mkdir() + sentinel = sibling / "keep-me" + sentinel.write_text("preserve", encoding="utf-8") + Path(index.file).symlink_to(sibling, target_is_directory=True) + + result = backend.build(_dataset(), [index], force=True) + + assert not result.success + assert "must not be a symlink" in result.error_message + assert Path(index.file).is_symlink() + assert sentinel.read_text(encoding="utf-8") == "preserve" + assert factory.calls == [] + + +def test_force_rebuild_removes_only_the_valid_index_directory( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + path = Path(index.file) + path.mkdir(parents=True) + stale_file = path / "stale" + stale_file.write_text("old", encoding="utf-8") + sibling = path.parent / "sibling" + sibling.mkdir() + sibling_file = sibling / "keep-me" + sibling_file.write_text("preserve", encoding="utf-8") + + result = backend.build(_dataset(), [index], force=True) + + assert result.success, result.error_message + assert not stale_file.exists() + assert sibling_file.read_text(encoding="utf-8") == "preserve" + + +def test_failed_force_rebuild_preserves_the_previous_valid_index( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + first_result = backend.build(_dataset(), [index]) + assert first_result.success, first_result.error_message + path = Path(index.file) + original_manifest = (path / ".cuvs-bench-lucene.json").read_bytes() + original_payload = (path / "segments.fake").read_bytes() + runtime.build_error = RuntimeError("replacement build failed") + + replacement = backend.build(_dataset(offset=0.25), [index], force=True) + + assert not replacement.success + assert ( + replacement.error_message == "RuntimeError: replacement build failed" + ) + assert (path / ".cuvs-bench-lucene.json").read_bytes() == original_manifest + assert (path / "segments.fake").read_bytes() == original_payload + assert list(path.parent.glob(f".{path.name}.build-*")) == [] + + +@pytest.mark.parametrize( + ("diagnostic", "message"), + ( + pytest.param( + "_index_size", "size measurement failed", id="index-size" + ), + pytest.param( + "_runtime_build_timing_metadata", + "runtime timing invalid", + id="runtime-timing", + ), + ), +) +def test_prepublication_diagnostic_failure_preserves_the_previous_index( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + diagnostic: str, + message: str, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + first_result = backend.build(_dataset(), [index]) + assert first_result.success, first_result.error_message + path = Path(index.file) + original_manifest = (path / ".cuvs-bench-lucene.json").read_bytes() + original_payload = (path / "segments.fake").read_bytes() + + def fail_diagnostic(_value: object) -> object: + raise OSError(message) + + monkeypatch.setattr( + f"cuvs_bench.backends.lucene.{diagnostic}", fail_diagnostic + ) + replacement = backend.build(_dataset(offset=0.25), [index], force=True) + + assert not replacement.success + assert replacement.error_message == f"OSError: {message}" + assert (path / ".cuvs-bench-lucene.json").read_bytes() == original_manifest + assert (path / "segments.fake").read_bytes() == original_payload + assert list(path.parent.glob(f".{path.name}.build-*")) == [] + + +def test_backup_cleanup_failure_does_not_report_a_published_index_as_failed( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + assert backend.build(_dataset(), [index]).success + + def fail_cleanup(_path: Path) -> None: + raise OSError("cleanup denied") + + monkeypatch.setattr( + "cuvs_bench.backends.lucene.shutil.rmtree", fail_cleanup + ) + replacement = backend.build(_dataset(offset=0.25), [index], force=True) + + assert replacement.success, replacement.error_message + assert "Published the new index" in replacement.metadata["cleanup_warning"] + assert "cleanup denied" in replacement.metadata["cleanup_warning"] + assert Path(index.file).is_dir() + + +def test_failed_index_install_restores_the_previous_index( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + destination = tmp_path / "index" + staged = tmp_path / "staged" + destination.mkdir() + staged.mkdir() + (destination / "old").write_text("preserve", encoding="utf-8") + original_rename = Path.rename + + def fail_install(path: Path, target: Path) -> Path: + if path == staged: + raise OSError("install denied") + return original_rename(path, target) + + monkeypatch.setattr(Path, "rename", fail_install) + + with pytest.raises(OSError, match="install denied"): + LuceneBackend._install_staged_index(staged, destination) + + assert (destination / "old").read_text(encoding="utf-8") == "preserve" + + +def test_failed_index_restore_identifies_the_recoverable_backup( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + destination = tmp_path / "index" + staged = tmp_path / "staged" + destination.mkdir() + staged.mkdir() + original_rename = Path.rename + + def fail_install_and_restore(path: Path, target: Path) -> Path: + if path == staged: + raise OSError("install denied") + if path.name.startswith(f".{destination.name}.backup-"): + raise OSError("restore denied") + return original_rename(path, target) + + monkeypatch.setattr(Path, "rename", fail_install_and_restore) + + with pytest.raises(OSError, match="install denied") as failure: + LuceneBackend._install_staged_index(staged, destination) + + [note] = failure.value.__notes__ + assert "Failed to restore the previous index" in note + assert "restore denied" in note + assert "recoverable backup" in note + + +@pytest.mark.parametrize("control_error", (KeyboardInterrupt, SystemExit)) +def test_restore_process_control_takes_precedence_over_install_failure( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + control_error: type[BaseException], +) -> None: + destination = tmp_path / "index" + staged = tmp_path / "staged" + destination.mkdir() + staged.mkdir() + original_rename = Path.rename + + def fail_install_then_interrupt_restore(path: Path, target: Path) -> Path: + if path == staged: + raise OSError("install denied") + if path.name.startswith(f".{destination.name}.backup-"): + raise control_error("restore interrupted") + return original_rename(path, target) + + monkeypatch.setattr(Path, "rename", fail_install_then_interrupt_restore) + + with pytest.raises(control_error) as failure: + LuceneBackend._install_staged_index(staged, destination) + + [note] = failure.value.__notes__ + assert "Index installation first failed: OSError: install denied" in note + assert "The previous index remains in" in note + assert list(tmp_path.glob(".index.backup-*")) + + +def test_backup_cleanup_does_not_swallow_process_control_exceptions( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + destination = tmp_path / "index" + staged = tmp_path / "staged" + destination.mkdir() + staged.mkdir() + + def interrupt_cleanup(_path: Path) -> None: + raise KeyboardInterrupt + + monkeypatch.setattr( + "cuvs_bench.backends.lucene.shutil.rmtree", interrupt_cleanup + ) + + with pytest.raises(KeyboardInterrupt): + LuceneBackend._install_staged_index(staged, destination) + + +def test_reusing_an_index_rejects_a_different_dataset_fingerprint( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + first_result = backend.build(_dataset(), [index]) + assert first_result.success, first_result.error_message + + reuse_result = backend.build(_dataset(offset=0.25), [index]) + + assert not reuse_result.success + assert ( + "does not match this dataset and configuration" + in reuse_result.error_message + ) + assert "rerun with --force" in reuse_result.error_message + assert len(runtime.build_calls) == 1 + + +def test_reusing_a_file_backed_index_does_not_materialize_training_vectors( + tmp_path: Path, +) -> None: + class ReuseOnlyDataset(Dataset): + @property + def training_vectors(self) -> np.ndarray: + raise AssertionError("index reuse materialized the base dataset") + + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + vectors = _dataset().training_vectors + base_file = tmp_path / "base.fbin" + _write_fbin(base_file, vectors) + build_dataset = Dataset( + name="tiny-l2", + training_vectors=vectors, + query_vectors=vectors[:2].copy(), + distance_metric="euclidean", + base_file=str(base_file), + ) + assert backend.build(build_dataset, [index]).success + reuse_dataset = ReuseOnlyDataset( + name="tiny-l2", + query_vectors=vectors[:2].copy(), + distance_metric="euclidean", + base_file=str(base_file), + ) + + result = backend.build(reuse_dataset, [index]) + + assert result.success, result.error_message + assert result.metadata["skipped"] is True + assert len(runtime.build_calls) == 1 + + +def test_search_rejects_changed_vectors_with_the_same_dataset_name( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + assert backend.build(_dataset(), [index]).success + + result = backend.search(_dataset(offset=0.25), [index], k=2)[0] + + assert not result.success + assert ( + "does not match this dataset and configuration" in result.error_message + ) + assert runtime.search_calls == [] + + +def test_file_backed_search_does_not_materialize_training_vectors( + tmp_path: Path, +) -> None: + class SearchOnlyDataset(Dataset): + @property + def training_vectors(self) -> np.ndarray: + raise AssertionError("search materialized the base dataset") + + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + vectors = _dataset().training_vectors + base_file = tmp_path / "base.fbin" + _write_fbin(base_file, vectors) + build_dataset = Dataset( + name="tiny-l2", + training_vectors=vectors, + query_vectors=vectors[:2].copy(), + distance_metric="euclidean", + base_file=str(base_file), + ) + assert backend.build(build_dataset, [index]).success + search_dataset = SearchOnlyDataset( + name="tiny-l2", + query_vectors=vectors[:2].copy(), + distance_metric="euclidean", + base_file=str(base_file), + ) + + result = backend.search(search_dataset, [index], k=2)[0] + + assert result.success, result.error_message + + +def test_file_backed_search_rejects_changed_content_when_tokens_collide( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + vectors = _dataset().training_vectors + base_file = tmp_path / "base.fbin" + _write_fbin(base_file, vectors) + dataset = Dataset( + name="tiny-l2", + training_vectors=vectors, + query_vectors=vectors[:2].copy(), + distance_metric="euclidean", + base_file=str(base_file), + ) + assert backend.build(dataset, [index]).success + original_source = _source_identity(dataset) + changed_vectors = vectors.copy() + changed_vectors[0, 0] = 42.0 + _write_fbin(base_file, changed_vectors) + monkeypatch.setattr( + "cuvs_bench.backends.lucene._source_identity", + lambda _dataset: original_source, + ) + search_dataset = Dataset( + name="tiny-l2", + query_vectors=vectors[:2].copy(), + distance_metric="euclidean", + base_file=str(base_file), + ) + + result = backend.search(search_dataset, [index], k=2)[0] + + assert not result.success + assert ( + "does not match this dataset and configuration" in result.error_message + ) + assert runtime.search_calls == [] + + +def test_search_rejects_a_file_that_was_not_the_explicit_build_array( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + indexed_vectors = _dataset().training_vectors + file_vectors = indexed_vectors.copy() + file_vectors[0, 0] = 42.0 + base_file = tmp_path / "base.fbin" + _write_fbin(base_file, file_vectors) + build_dataset = Dataset( + name="tiny-l2", + training_vectors=indexed_vectors, + query_vectors=indexed_vectors[:2].copy(), + distance_metric="euclidean", + base_file=str(base_file), + ) + assert backend.build(build_dataset, [index]).success + search_dataset = Dataset( + name="tiny-l2", + query_vectors=indexed_vectors[:2].copy(), + distance_metric="euclidean", + base_file=str(base_file), + ) + + result = backend.search(search_dataset, [index], k=2)[0] + + assert not result.success + assert ( + "does not match this dataset and configuration" in result.error_message + ) + assert runtime.search_calls == [] + + +def test_search_rejects_a_physical_index_that_fails_verification( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = _dataset() + assert backend.build(dataset, [index]).success + runtime.verification_error = RuntimeError( + "Segment '_0' uses CuVS2510GPUSearchCodec, not Lucene101" + ) + + result = backend.search(dataset, [index], k=2)[0] + + assert not result.success + assert "uses CuVS2510GPUSearchCodec, not Lucene101" in result.error_message + assert runtime.search_calls == [] + + +def test_search_rejects_malformed_build_runtime_provenance( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CAGRA_ALGORITHM, runtime + ) + dataset = _dataset() + assert backend.build(dataset, [index]).success + manifest_path = Path(index.file) / ".cuvs-bench-lucene.json" + manifest = json.loads(manifest_path.read_text(encoding="utf-8")) + manifest["build_runtime_artifacts"]["cuvs_lucene_jar_sha256"] = "invalid" + manifest_path.write_text(json.dumps(manifest), encoding="utf-8") + + result = backend.search(dataset, [index], k=2)[0] + + assert not result.success + assert "invalid cuvs_lucene_jar_sha256" in result.error_message + assert runtime.search_calls == [] + + +def test_search_requests_rebuild_for_the_legacy_stored_id_schema( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = _dataset() + assert backend.build(dataset, [index]).success + manifest_path = Path(index.file) / ".cuvs-bench-lucene.json" + manifest = json.loads(manifest_path.read_text(encoding="utf-8")) + manifest["schema_version"] = 1 + manifest_path.write_text(json.dumps(manifest), encoding="utf-8") + + result = backend.search(dataset, [index], k=2)[0] + + assert not result.success + assert "Unsupported Lucene index manifest" in result.error_message + assert "rerun with --force" in result.error_message + assert runtime.search_calls == [] + + +@pytest.mark.parametrize( + ("field", "value", "message"), + ( + pytest.param( + "schema_version", + True, + "Unsupported Lucene index manifest", + id="boolean-schema-version", + ), + pytest.param( + "segment_count", + 1.0, + "invalid segment count", + id="floating-segment-count", + ), + pytest.param( + "segment_count", + 0, + "invalid segment count", + id="empty-index-segment-count", + ), + ), +) +def test_search_rejects_invalid_manifest_integer_fields( + tmp_path: Path, field: str, value: object, message: str +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = _dataset() + assert backend.build(dataset, [index]).success + manifest_path = Path(index.file) / ".cuvs-bench-lucene.json" + manifest = json.loads(manifest_path.read_text(encoding="utf-8")) + manifest[field] = value + manifest_path.write_text(json.dumps(manifest), encoding="utf-8") + + result = backend.search(dataset, [index], k=2)[0] + + assert not result.success + assert message in result.error_message + assert runtime.search_calls == [] + + +def test_search_rejects_non_integer_manifest_dataset_dimensions( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = _dataset() + assert backend.build(dataset, [index]).success + manifest_path = Path(index.file) / ".cuvs-bench-lucene.json" + manifest = json.loads(manifest_path.read_text(encoding="utf-8")) + manifest["dataset"]["dimensions"] = 2.0 + manifest_path.write_text(json.dumps(manifest), encoding="utf-8") + + result = backend.search(dataset, [index], k=2)[0] + + assert not result.success + assert "invalid dataset dimensions" in result.error_message + assert runtime.search_calls == [] + + +def test_search_rejects_manifest_segment_count_that_disagrees_with_index( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + dataset = _dataset() + assert backend.build(dataset, [index]).success + runtime.segment_count = 2 + + result = backend.search(dataset, [index], k=2)[0] + + assert not result.success + assert "segment count does not match the physical index: 1 != 2" in ( + result.error_message + ) + assert runtime.search_calls == [] + + +def test_reusing_the_same_index_reports_a_skipped_build( + tmp_path: Path, +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index( + tmp_path, CPU_HNSW_ALGORITHM, runtime + ) + assert backend.build(_dataset(), [index]).success + assert runtime.artifact_verification_count == 1 + + reuse_result = backend.build(_dataset(), [index]) + + assert reuse_result.success, reuse_result.error_message + assert reuse_result.metadata == { + "skipped": True, + "codec": CPU_HNSW_CODEC, + "group": "test", + "index_name": CPU_HNSW_ALGORITHM, + "persisted_index_kind": "cpu_hnsw", + "build_route_policy": "cpu_hnsw", + "segment_count": 1, + "field_count": 1, + "vector_count": 4, + "dimensions": 2, + } + assert len(runtime.build_calls) == 1 + assert runtime.artifact_verification_count == 2 diff --git a/python/cuvs_bench/cuvs_bench/tests/test_lucene_results_config.py b/python/cuvs_bench/cuvs_bench/tests/test_lucene_results_config.py new file mode 100644 index 0000000000..4ae3837fe7 --- /dev/null +++ b/python/cuvs_bench/cuvs_bench/tests/test_lucene_results_config.py @@ -0,0 +1,293 @@ +# +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# + +"""Result, parameter, and configuration contracts for the Lucene backend.""" + +from __future__ import annotations + +from pathlib import Path +from typing import Any + +import numpy as np +import pytest + +from _lucene_test_support import ( + ALGORITHM_CASES, + RecordingRuntime, + _backend_and_index, +) +from cuvs_bench.backends._lucene_runtime import ( + ACCELERATED_HNSW_CODEC, + CAGRA_CODEC, + CPU_HNSW_CODEC, +) +from cuvs_bench.backends.lucene import ( + ACCELERATED_HNSW_ALGORITHM, + CAGRA_ALGORITHM, + CPU_HNSW_ALGORITHM, + LuceneBackend, + LuceneConfigLoader, + _codec_for, + _score_to_squared_euclidean, + _search_parameters, +) + + +@pytest.mark.parametrize(("algorithm", "codec", "_path"), ALGORITHM_CASES) +def test_backend_accepts_only_the_codec_owned_by_each_algorithm( + tmp_path: Path, algorithm: str, codec: str, _path: str +) -> None: + runtime = RecordingRuntime() + backend, index, _factory = _backend_and_index(tmp_path, algorithm, runtime) + + assert backend.algorithm == algorithm + assert backend.codec == codec + assert _codec_for(index.algo, index.build_param) == codec + + +@pytest.mark.parametrize( + ("algorithm", "codec"), + ( + pytest.param(CPU_HNSW_ALGORITHM, CAGRA_CODEC, id="cpu-with-cagra"), + pytest.param(CAGRA_ALGORITHM, CPU_HNSW_CODEC, id="cagra-with-cpu"), + pytest.param( + ACCELERATED_HNSW_ALGORITHM, + CPU_HNSW_CODEC, + id="accelerated-with-cpu", + ), + ), +) +def test_backend_rejects_mismatched_algorithm_and_codec( + tmp_path: Path, algorithm: str, codec: str +) -> None: + with pytest.raises( + ValueError, match="Invalid Lucene algorithm/codec pair" + ): + LuceneBackend( + { + "name": "invalid", + "algo": algorithm, + "codec": codec, + "index_root": str(tmp_path), + } + ) + + +@pytest.mark.parametrize(("algorithm", "codec", "_path"), ALGORITHM_CASES) +def test_build_parameters_accept_only_the_fixed_codec( + algorithm: str, codec: str, _path: str +) -> None: + assert _codec_for(algorithm, {}) == codec + assert _codec_for(algorithm, {"codec": codec}) == codec + + with pytest.raises( + ValueError, match="Unsupported Lucene build parameters" + ): + _codec_for(algorithm, {"codec": codec, "graph_degree": 32}) + + +@pytest.mark.parametrize( + ("parameters", "k", "expected"), + ( + pytest.param({}, 3, {"num_candidates": 3}, id="default-to-k"), + pytest.param( + {"num_candidates": 7}, + 3, + {"num_candidates": 7}, + id="explicit-candidates", + ), + ), +) +@pytest.mark.parametrize( + "algorithm", + (CPU_HNSW_ALGORITHM, ACCELERATED_HNSW_ALGORITHM), +) +def test_hnsw_search_parameter_validation_accepts_candidate_budgets( + algorithm: str, + parameters: dict[str, int], + k: int, + expected: dict[str, int], +) -> None: + assert _search_parameters(algorithm, parameters, k) == expected + + +@pytest.mark.parametrize( + "parameters", + ( + pytest.param({"num_candidates": 2}, id="below-k"), + pytest.param({"num_candidates": True}, id="boolean"), + pytest.param({"search_width": 16}, id="unsupported-key"), + ), +) +@pytest.mark.parametrize( + "algorithm", + (CPU_HNSW_ALGORITHM, ACCELERATED_HNSW_ALGORITHM), +) +def test_hnsw_search_parameter_validation_rejects_invalid_budgets( + algorithm: str, + parameters: dict[str, Any], +) -> None: + with pytest.raises(ValueError): + _search_parameters(algorithm, parameters, 3) + + +def test_cagra_search_accepts_only_fixed_parameters_within_its_route_limit() -> ( + None +): + assert _search_parameters(CAGRA_ALGORITHM, {}, 1024) == { + "num_candidates": 1024 + } + + with pytest.raises(ValueError, match="accepts no search parameters"): + _search_parameters(CAGRA_ALGORITHM, {"search_width": 16}, 10) + with pytest.raises(ValueError, match=r"supports k <= 1024"): + _search_parameters(CAGRA_ALGORITHM, {}, 1025) + + +@pytest.mark.parametrize( + ("score", "expected_distance"), + ( + pytest.param(1.0, 0.0, id="zero-distance"), + pytest.param(0.5, 1.0, id="unit-distance"), + pytest.param(0.2, 4.0, id="distance-four"), + ), +) +def test_lucene_scores_are_inverted_to_squared_euclidean_distance( + score: float, expected_distance: float +) -> None: + assert _score_to_squared_euclidean(score) == pytest.approx( + expected_distance + ) + + +@pytest.mark.parametrize( + "score", + ( + 0.0, + -1.0, + 1.0 + 2.0 * float(np.spacing(np.float32(1.0))), + np.nan, + np.inf, + ), +) +def test_score_inversion_rejects_values_outside_lucenes_score_domain( + score: float, +) -> None: + with pytest.raises(RuntimeError, match="invalid Euclidean score"): + _score_to_squared_euclidean(score) + + +def test_score_inversion_tolerates_float32_roundoff_above_one() -> None: + score = float(np.nextafter(np.float32(1.0), np.float32(2.0))) + + assert _score_to_squared_euclidean(score) == 0.0 + + +def test_config_loader_maps_each_algorithm_to_its_codec_and_requirement( + tmp_path: Path, +) -> None: + dataset_configuration = tmp_path / "datasets.yaml" + dataset_configuration.write_text( + "- name: tiny-l2\n distance: euclidean\n dims: 2\n", + encoding="utf-8", + ) + loader = LuceneConfigLoader() + + _dataset_config, configurations = loader.load( + dataset="tiny-l2", + dataset_path=str(tmp_path), + dataset_configuration=str(dataset_configuration), + algorithms=( + f"{CPU_HNSW_ALGORITHM},{ACCELERATED_HNSW_ALGORITHM}," + f"{CAGRA_ALGORITHM}" + ), + groups="test", + ) + by_algorithm = { + configuration.indexes[0].algo: configuration + for configuration in configurations + } + + assert set(by_algorithm) == { + CPU_HNSW_ALGORITHM, + ACCELERATED_HNSW_ALGORITHM, + CAGRA_ALGORITHM, + } + assert by_algorithm[CPU_HNSW_ALGORITHM].indexes[0].build_param == { + "codec": CPU_HNSW_CODEC + } + assert by_algorithm[CAGRA_ALGORITHM].indexes[0].build_param == { + "codec": CAGRA_CODEC + } + assert by_algorithm[ACCELERATED_HNSW_ALGORITHM].indexes[0].build_param == { + "codec": ACCELERATED_HNSW_CODEC + } + assert ( + by_algorithm[CPU_HNSW_ALGORITHM].backend_config["requires_cuvs"] + is False + ) + assert ( + by_algorithm[ACCELERATED_HNSW_ALGORITHM].backend_config[ + "requires_cuvs" + ] + is True + ) + assert ( + by_algorithm[CAGRA_ALGORITHM].backend_config["requires_cuvs"] is True + ) + assert ( + by_algorithm[CPU_HNSW_ALGORITHM].backend_config["include_cuvs"] is True + ) + assert ( + by_algorithm[ACCELERATED_HNSW_ALGORITHM].backend_config["include_cuvs"] + is True + ) + assert by_algorithm[CAGRA_ALGORITHM].backend_config["include_cuvs"] is True + assert by_algorithm[CPU_HNSW_ALGORITHM].backend_config["group"] == "test" + assert ( + by_algorithm[ACCELERATED_HNSW_ALGORITHM].backend_config["group"] + == "test" + ) + assert by_algorithm[CAGRA_ALGORITHM].backend_config["group"] == "test" + + +def test_config_loader_rejects_dataset_names_that_can_escape_the_root( + tmp_path: Path, +) -> None: + dataset_configuration = tmp_path / "datasets.yaml" + dataset_configuration.write_text( + "- name: ../escape\n distance: euclidean\n dims: 2\n", + encoding="utf-8", + ) + + with pytest.raises(ValueError, match="Unsafe Lucene dataset"): + LuceneConfigLoader().load( + dataset="../escape", + dataset_path=str(tmp_path), + dataset_configuration=str(dataset_configuration), + algorithms=CPU_HNSW_ALGORITHM, + groups="test", + ) + + +def test_cpu_only_config_does_not_include_optional_cuvs_artifacts( + tmp_path: Path, +) -> None: + dataset_configuration = tmp_path / "datasets.yaml" + dataset_configuration.write_text( + "- name: tiny-l2\n distance: euclidean\n dims: 2\n", + encoding="utf-8", + ) + + _dataset_config, [configuration] = LuceneConfigLoader().load( + dataset="tiny-l2", + dataset_path=str(tmp_path), + dataset_configuration=str(dataset_configuration), + algorithms=CPU_HNSW_ALGORITHM, + groups="test", + ) + + assert configuration.backend_config["requires_cuvs"] is False + assert configuration.backend_config["include_cuvs"] is False From 4f7d192f7e9fbe8fad5702441db616c741bbbaf3 Mon Sep 17 00:00:00 2001 From: nvzm123 Date: Tue, 6 Oct 2026 20:58:00 +0000 Subject: [PATCH 2/2] Preserve Lucene benchmark groups --- .../cuvs_bench/cuvs_bench/backends/lucene.py | 6 +- .../tests/test_lucene_results_config.py | 74 +++++++++++++++++++ 2 files changed, 79 insertions(+), 1 deletion(-) diff --git a/python/cuvs_bench/cuvs_bench/backends/lucene.py b/python/cuvs_bench/cuvs_bench/backends/lucene.py index 77d728b429..dfa2473a67 100644 --- a/python/cuvs_bench/cuvs_bench/backends/lucene.py +++ b/python/cuvs_bench/cuvs_bench/backends/lucene.py @@ -778,7 +778,11 @@ def _discover_algo_groups( config = configs.get(algorithm) if config is None: raise ValueError(f"No configuration found for {algorithm!r}") - selected_groups = requested_pairs.get(algorithm, groups) + selected_groups = list( + dict.fromkeys( + [*groups, *requested_pairs.get(algorithm, [])] + ) + ) for group in selected_groups: try: group_config = config["groups"][group] diff --git a/python/cuvs_bench/cuvs_bench/tests/test_lucene_results_config.py b/python/cuvs_bench/cuvs_bench/tests/test_lucene_results_config.py index 4ae3837fe7..2cdffabc9c 100644 --- a/python/cuvs_bench/cuvs_bench/tests/test_lucene_results_config.py +++ b/python/cuvs_bench/cuvs_bench/tests/test_lucene_results_config.py @@ -253,6 +253,80 @@ def test_config_loader_maps_each_algorithm_to_its_codec_and_requirement( assert by_algorithm[CAGRA_ALGORITHM].backend_config["group"] == "test" +@pytest.mark.parametrize( + ("groups", "algorithm_groups", "expected_groups"), + ( + pytest.param("base", None, ["base"], id="common-only"), + pytest.param( + "base", + f"{CPU_HNSW_ALGORITHM}.test", + ["base", "test"], + id="add-algorithm-group", + ), + pytest.param( + "test", + f"{CPU_HNSW_ALGORITHM}.base", + ["test", "base"], + id="preserve-order", + ), + pytest.param( + "base,test", + ( + f"{CPU_HNSW_ALGORITHM}.test," + f"{CPU_HNSW_ALGORITHM}.base" + ), + ["base", "test"], + id="deduplicate", + ), + ), +) +def test_config_loader_adds_algorithm_groups_to_common_groups( + tmp_path: Path, + groups: str, + algorithm_groups: str | None, + expected_groups: list[str], +) -> None: + dataset_configuration = tmp_path / "datasets.yaml" + dataset_configuration.write_text( + "- name: tiny-l2\n distance: euclidean\n dims: 2\n", + encoding="utf-8", + ) + + _dataset_config, configurations = LuceneConfigLoader().load( + dataset="tiny-l2", + dataset_path=str(tmp_path), + dataset_configuration=str(dataset_configuration), + algorithms=CPU_HNSW_ALGORITHM, + groups=groups, + algo_groups=algorithm_groups, + ) + + assert [ + configuration.backend_config["group"] + for configuration in configurations + ] == expected_groups + + +def test_config_loader_rejects_unknown_algorithm_group( + tmp_path: Path, +) -> None: + dataset_configuration = tmp_path / "datasets.yaml" + dataset_configuration.write_text( + "- name: tiny-l2\n distance: euclidean\n dims: 2\n", + encoding="utf-8", + ) + + with pytest.raises(ValueError, match="No Lucene group 'missing'"): + LuceneConfigLoader().load( + dataset="tiny-l2", + dataset_path=str(tmp_path), + dataset_configuration=str(dataset_configuration), + algorithms=CPU_HNSW_ALGORITHM, + groups="base", + algo_groups=f"{CPU_HNSW_ALGORITHM}.missing", + ) + + def test_config_loader_rejects_dataset_names_that_can_escape_the_root( tmp_path: Path, ) -> None: