diff --git a/nsparse/CMakeLists.txt b/nsparse/CMakeLists.txt index 3aa6bd1..27cb615 100644 --- a/nsparse/CMakeLists.txt +++ b/nsparse/CMakeLists.txt @@ -11,6 +11,9 @@ set(NSPARSE_SRC seismic_index.cpp seismic_scalar_quantized_index.cpp disk_seismic_index.cpp + disk_seismic_scalar_quantized_index.cpp + disk_seismic_index_base.cpp + disk_seismic_search.cpp sparse_vectors.cpp index_factory.cpp id_map_index.cpp @@ -34,6 +37,10 @@ set(NSPARSE_HEADERS mmap_index.h seismic_index.h seismic_scalar_quantized_index.h + disk_seismic_index.h + disk_seismic_scalar_quantized_index.h + disk_seismic_index_base.h + disk_seismic_search.h inline_forward_index.h index_factory.h id_map_index.h diff --git a/nsparse/cluster/inverted_list_clusters.h b/nsparse/cluster/inverted_list_clusters.h index 8839c11..9d89bda 100644 --- a/nsparse/cluster/inverted_list_clusters.h +++ b/nsparse/cluster/inverted_list_clusters.h @@ -39,6 +39,11 @@ class InvertedListClusters : public MmapSerializable { size_t cluster_size() const { return n_clusters_; } + // Width (bytes) of each stored summary value: the dispatch + // score_summaries_transposed uses, so a reader can check it against the + // width its query is encoded at. + [[nodiscard]] size_t element_size() const { return element_size_; } + void serialize(IOWriter* writer) const override; // Like serialize() but emits the doc-id membership empty, for indexes that // keep it elsewhere (DiskSeismic's inline forward index). Read back by the diff --git a/nsparse/disk_seismic_index.cpp b/nsparse/disk_seismic_index.cpp index 4cc6cb4..e375cc4 100644 --- a/nsparse/disk_seismic_index.cpp +++ b/nsparse/disk_seismic_index.cpp @@ -9,317 +9,42 @@ #include "nsparse/disk_seismic_index.h" -#include #include #include #include -#include #include #include -#include "absl/container/flat_hash_set.h" -#include "nsparse/cluster/inverted_list_clusters.h" -#include "nsparse/id_selector.h" -#include "nsparse/index.h" -#include "nsparse/io/inline_forward_index_io.h" -#include "nsparse/io/seismic_invlists_writer.h" +#include "nsparse/disk_seismic_index_base.h" #include "nsparse/seismic_common.h" -#include "nsparse/sparse_vectors.h" #include "nsparse/types.h" #include "nsparse/utils/checks.h" -#include "nsparse/utils/distance_simd.h" #include "nsparse/utils/mmap_cursor.h" #include "nsparse/utils/mmap_file.h" -#include "nsparse/utils/ranker.h" -#include "nsparse/utils/vector_process.h" namespace nsparse { -namespace { - -constexpr int kElementSize = U32; - -// Score one doc's vector against the dense query and push it, honoring -// visited-dedup and the id selector. Shared by both vector sources. -inline void score_doc(idx_t doc_id, const term_t* comps, const float* vals, - size_t nnz, const std::vector& dense, - const IDSelector* id_selector, - detail::TopKHolder& heap, - absl::flat_hash_set& visited) { - if (!visited.insert(doc_id).second) { - return; - } - if (id_selector != nullptr && !id_selector->is_member(doc_id)) { - return; - } - const float score = - detail::dot_product_float_dense(comps, vals, nnz, dense.data()); - heap.add(score, doc_id); -} - -// A candidate block: its summary score and its (posting list, cluster) address. -struct BlockCandidate { - float score; - uint32_t pl; - uint32_t cid; -}; - -// Score every document of one block against the dense query, dedup via -// `visited` and honor the id selector. Vectors come from the inline forward -// index (fwd) when loaded, else from the in-RAM CSR of a fresh build. -void score_block(const detail::InlineForwardIndex* fwd, - const SparseVectors* vectors, - const std::vector& clusters, uint32_t pl, - uint32_t cid, const std::vector& dense, - const IDSelector* id_selector, size_t element_size, - detail::TopKHolder& heap, - absl::flat_hash_set& visited) { - if (fwd != nullptr) { - const detail::BlockView bv = fwd->block(pl, cid); - if (bv.absent()) { - return; - } - for (uint32_t i = 0; i < bv.n_docs; ++i) { - score_doc( - static_cast(bv.doc_ids[i]), bv.doc_comps(i), - reinterpret_cast(bv.doc_vals(i, element_size)), - bv.nnz(i), dense, id_selector, heap, visited); - } - } else if (vectors != nullptr) { - const auto data = vectors->get_all_data(); - const idx_t* const indptr = data.indptr_data; - const term_t* const indices = data.indices_data; - const float* const values = data.values_data; - for (const idx_t doc_id : clusters[pl].get_docs(cid)) { - const idx_t start = indptr[doc_id]; - const size_t len = indptr[doc_id + 1] - start; - score_doc(doc_id, indices + start, values + start, len, dense, - id_selector, heap, visited); - } - } -} -} // namespace DiskSeismicIndex::DiskSeismicIndex(int dim) - : MmapIndex(dim), - cluster_parameter_(detail::kDefaultSeismicClusterParams) {} + : DiskSeismicIndexBase(dim, detail::kDefaultSeismicClusterParams) {} DiskSeismicIndex::DiskSeismicIndex(int dim, SeismicClusterParameters parameter) - : MmapIndex(dim), cluster_parameter_(parameter) {} - -void DiskSeismicIndex::add(idx_t n, const idx_t* indptr, const term_t* indices, - const float* values) { - throw_if_not_positive(n); - throw_if_any_null(indptr, indices, values); - const size_t indptr_size = n + 1; - const size_t nnz = indptr[n]; - if (vectors_ == nullptr) { - // Fresh container: start the count at 0 so a stale num_vectors_ (e.g. - // left by a prior mmap load, which has no vectors_) cannot accumulate. - num_vectors_ = 0; - vectors_ = std::make_unique( - SparseVectorsConfig{.element_size = kElementSize, - .dimension = static_cast(dimension_)}); - } - vectors_->add_vectors(indptr, indptr_size, indices, nnz, - reinterpret_cast(values), - nnz * kElementSize); - num_vectors_ += n; -} - -void DiskSeismicIndex::build() { - clustered_inverted_lists = detail::build_inverted_lists_clusters( - get_vectors(), - {.element_size = kElementSize, - .dimension = static_cast(get_dimension())}, - cluster_parameter_); -} - -auto DiskSeismicIndex::search(idx_t n, const idx_t* indptr, - const term_t* indices, const float* values, int k, - SearchParameters* search_parameters) - -> pair_of_score_id_vectors_t { - // Quit early when there is nothing to score: no vectors, no queries, or no - // forward-vector source (fwd_ empty and vectors_ null — a corrupt or - // uninitialized index). - if (num_vectors_ == 0 || n == 0 || - (fwd_.num_blocks() == 0 && vectors_ == nullptr)) { - return { - std::vector>(n, std::vector(k, -1.0F)), - std::vector>( - n, std::vector(k, detail::INVALID_IDX))}; - } - - // Resolve `cut` and `k_prime`. A DiskSeismicSearchParameters carries - // k_prime; a plain SeismicSearchParameters (or null) uses the default - // budget. The id-selector exact-match fast path is omitted (it needs an - // in-RAM SparseVectors a mapped index lacks); the selector is still honored - // in score_block. - const DiskSeismicSearchParameters defaults; - int cut = defaults.cut; - int k_prime = defaults.k_prime; - if (const auto* disk_parameters = - dynamic_cast( - search_parameters)) { - cut = disk_parameters->cut; - k_prime = disk_parameters->k_prime; - } else if (const auto* seismic_parameters = - dynamic_cast( - search_parameters)) { - cut = seismic_parameters->cut; - } - // else: search_parameters is null or an unrelated SearchParameters subtype, - // so cut and k_prime keep the defaults initialized above. + : DiskSeismicIndexBase(dim, parameter) {} - // k_prime is a block budget, not a document count: a block holds many docs, - // so k_prime < k is valid (a few blocks can still fill k, and a short-fall - // is padded like any under-budget search). Only a non-positive budget is - // rejected. - if (k_prime <= 0) { - throw std::invalid_argument( - "DiskSeismicIndex: k_prime (block budget) must be positive"); - } +size_t DiskSeismicIndex::code_element_size() const { return U32; } - std::vector> result_distances(n, std::vector(k, -1.0F)); - std::vector> result_labels( - n, std::vector(k, detail::INVALID_IDX)); - const size_t dim = static_cast(dimension_); - -#pragma omp parallel - { - std::vector dense(dim, 0.0F); - absl::flat_hash_set visited; - visited.reserve(static_cast(std::max(k, 1)) * 4096); - -#pragma omp for schedule(dynamic, 64) - for (idx_t query_idx = 0; query_idx < n; ++query_idx) { - const idx_t start = indptr[query_idx]; - const size_t len = indptr[query_idx + 1] - start; - const term_t* q_indices = indices + start; - const float* q_values = values + start; - const auto& cuts = - detail::top_k_tokens(q_indices, q_values, len, cut); - auto [distances, labels] = - single_query(dense, visited, q_indices, q_values, len, cuts, k, - k_prime, search_parameters); - result_distances[query_idx] = std::move(distances); - result_labels[query_idx] = std::move(labels); - } - } - return {result_distances, result_labels}; -} - -auto DiskSeismicIndex::single_query( - std::vector& dense, absl::flat_hash_set& visited, - const term_t* q_indices, const float* q_values, size_t q_len, - const std::vector& cuts, int k, int k_prime, - SearchParameters* search_parameters) -> pair_of_score_id_vector_t { - if (num_vectors_ == 0) { - return {{}, {}}; - } - // Scatter the query into the dense lookup table. - for (size_t i = 0; i < q_len; ++i) { - dense[q_indices[i]] = q_values[i]; - } - visited.clear(); - - // Prefer the mapped inline forward index; fall back to the in-RAM vectors - // of a fresh build. - const detail::InlineForwardIndex* fwd = - fwd_.num_blocks() > 0 ? &fwd_ : nullptr; - const SparseVectors* vectors = fwd == nullptr ? vectors_.get() : nullptr; - const size_t element_size = fwd != nullptr - ? fwd->element_size() - : static_cast(kElementSize); - const IDSelector* id_selector = search_parameters == nullptr - ? nullptr - : search_parameters->get_id_selector(); - - // Collect every candidate block across the cut lists with its summary - // score. score_summaries_transposed dots the full query with each cluster - // summary, so scores are comparable across posting lists for the global - // ranking. - std::vector candidates; - std::vector score_scratch; - for (const term_t term : cuts) { - if (term >= clustered_inverted_lists.size()) [[unlikely]] { - continue; - } - const InvertedListClusters& cluster_invlist = - clustered_inverted_lists[term]; - const size_t n_clusters = cluster_invlist.cluster_size(); - if (n_clusters == 0) { - continue; - } - cluster_invlist.score_summaries_transposed( - q_indices, reinterpret_cast(q_values), q_len, - score_scratch); - for (size_t cid = 0; cid < n_clusters; ++cid) { - candidates.push_back( - {score_scratch[cid], term, static_cast(cid)}); - } - } - - // Select the global top-k_prime blocks by summary score. The order among - // them does not affect the result (dedup + exact scoring are - // order-independent). - const size_t budget = - std::min(static_cast(k_prime), candidates.size()); - if (budget < candidates.size()) { - std::nth_element(candidates.begin(), candidates.begin() + budget, - candidates.end(), - [](const BlockCandidate& a, const BlockCandidate& b) { - return a.score > b.score; - }); - candidates.resize(budget); - } - - // Score the selected blocks; visited dedups docs shared across them. - detail::TopKHolder holder(k); - for (const BlockCandidate& candidate : candidates) { - score_block(fwd, vectors, clustered_inverted_lists, candidate.pl, - candidate.cid, dense, id_selector, element_size, holder, - visited); - } - - // Reset only the query's own positions (q_len is small) instead of - // memset-ing the whole dim-sized table; mirrors the scatter at entry. The - // scattered writes are bounded by q_len and negligible next to the block - // reads/scoring that dominate the query. - for (size_t i = 0; i < q_len; ++i) { - dense[q_indices[i]] = 0.0F; - } - auto [scores, ids] = holder.top_k_items_descending(); - scores.resize(k, -1.0F); - ids.resize(k, detail::INVALID_IDX); - return {scores, ids}; +// Float values are stored verbatim, so their bytes are the codes -- no copy, +// `scratch` unused. +const uint8_t* DiskSeismicIndex::encode_values( + const float* values, size_t /*nnz*/, + std::vector& /*scratch*/) const { + return reinterpret_cast(values); } -void DiskSeismicIndex::write_index(IOWriter* io_writer) { - const uint64_t nv = num_vectors_; - io_writer->write(const_cast(&nv), sizeof(uint64_t), 1); - // Summaries only: the doc-id membership is already in the inline forward - // index below, so writing it in the posting lists too would duplicate it. - SeismicInvertedListsWriter inv_list_writer(clustered_inverted_lists, - /*summaries_only=*/true); - inv_list_writer.serialize(io_writer); - // Inline forward index, built from the same clusters + vectors. An empty - // corpus uses a correctly-typed empty SparseVectors (element_size must be a - // valid width even with zero vectors) so the section still round-trips. - SparseVectors empty_vectors({.element_size = kElementSize, - .dimension = static_cast(dimension_)}); - const SparseVectors& v = vectors_ != nullptr ? *vectors_ : empty_vectors; - detail::InlineForwardIndex forward(clustered_inverted_lists, v); - forward.serialize(io_writer); -} - -void DiskSeismicIndex::read_index(IOReader* /*io_reader*/, - const IndexHeader& /*header*/, - int /*io_flags*/) { - // The inline forward index is borrowed from a mapping, never copied onto - // the heap, so this index has no copying read path. - throw std::runtime_error( - "DiskSeismicIndex is mmap-only; load with read_index(file, " - "IndexIoFlag::kUseMmap)"); +const uint8_t* DiskSeismicIndex::encode_query( + const float* values, size_t /*nnz*/, + const SearchParameters* /*search_parameters*/, + std::vector& /*scratch*/) const { + return reinterpret_cast(values); } DiskSeismicIndex* DiskSeismicIndex::mmap_index(const IndexHeader& header, @@ -332,20 +57,8 @@ DiskSeismicIndex* DiskSeismicIndex::mmap_index(const IndexHeader& header, MmapCursor cursor(mmap_file.data(), mmap_file.size()); cursor.skip(pos); - // Same order write_index wrote them: doc count, summaries, inline forward. - const uint64_t nv = cursor.read_scalar(); - SeismicInvertedListsWriter inv_list_writer; - inv_list_writer.mmap_deserialize(&cursor); - detail::InlineForwardIndex forward; - forward.mmap_deserialize(&cursor); - - // Committed only once everything parsed. mapped_file_ last: the summaries - // and the forward index borrow from it, and moving it does not move the - // mapping. - index->clustered_inverted_lists = std::move(inv_list_writer.release()); - index->fwd_ = std::move(forward); - index->num_vectors_ = nv; - index->mapped_file_ = std::move(mmap_file); + // No extra header for the float index; the shared payload follows directly. + index->load_mapped_payload(&cursor, std::move(mmap_file)); return index.release(); } } // namespace nsparse diff --git a/nsparse/disk_seismic_index.h b/nsparse/disk_seismic_index.h index a5c2054..7ab1abf 100644 --- a/nsparse/disk_seismic_index.h +++ b/nsparse/disk_seismic_index.h @@ -14,12 +14,8 @@ #include #include -#include "absl/container/flat_hash_set.h" -#include "nsparse/cluster/inverted_list_clusters.h" -#include "nsparse/io/inline_forward_index_io.h" +#include "nsparse/disk_seismic_index_base.h" #include "nsparse/io/io.h" -#include "nsparse/mmap_index.h" -#include "nsparse/seismic_common.h" #include "nsparse/seismic_index.h" // SeismicSearchParameters #include "nsparse/types.h" @@ -30,7 +26,7 @@ inline constexpr int kDefaultBlockBudget = 50; // `cut` (inherited) bounds which posting lists' summaries are scored; `k_prime` // caps how many blocks are read (the k_prime highest-scoring ones). Replaces -// the inherited heap_factor, which DiskSeismicIndex ignores. +// the inherited heap_factor, which the disk-resident indexes ignore. struct DiskSeismicSearchParameters : public SeismicSearchParameters { int k_prime = kDefaultBlockBudget; DiskSeismicSearchParameters() = default; @@ -39,18 +35,13 @@ struct DiskSeismicSearchParameters : public SeismicSearchParameters { k_prime(k_prime) {} }; -// A seismic index whose per-document forward vectors live on disk in the -// block-contiguous (inline) layout of io/inline_forward_index_io.h, borrowed -// via mmap at search time; the cluster summaries stay in RAM. Search scores -// every candidate block's summary across the `cut` posting lists, selects the -// global top-k_prime blocks, and reads and scores only those. Pass a -// DiskSeismicSearchParameters to set k_prime, else kDefaultBlockBudget is used. -// -// A block's vectors come from fwd_.block() (the mapping) once loaded, else from -// the in-RAM SparseVectors of a fresh build. +// A SEISMIC index whose per-document forward vectors live on disk as float, in +// the block-contiguous (inline) layout, borrowed via mmap at search time; the +// cluster summaries stay in RAM. The disk-resident search and serialization +// live in DiskSeismicIndexBase; this type only pins the value width to float. // // mmap-only: load with read_index(file, kUseMmap); the copying read throws. -class DiskSeismicIndex : public MmapIndex, public IndexIO { +class DiskSeismicIndex : public DiskSeismicIndexBase { public: static constexpr std::array name = {'D', 'S', 'E', 'I'}; // Bump whenever write_index's payload layout changes. @@ -64,51 +55,22 @@ class DiskSeismicIndex : public MmapIndex, public IndexIO { DiskSeismicIndex(const DiskSeismicIndex&) = delete; DiskSeismicIndex& operator=(const DiskSeismicIndex&) = delete; - void add(idx_t n, const idx_t* indptr, const term_t* indices, - const float* values) override; - void build() override; - - // Persisted, since a mapped index has no in-RAM vectors_ to derive it from. - size_t num_vectors() const override { return num_vectors_; } - // Borrows a serialized index from a file mapping. `pos` is where the // payload begins. static DiskSeismicIndex* mmap_index(const IndexHeader& header, const char* index_file, size_t pos); -protected: - std::vector clustered_inverted_lists; - private: [[nodiscard]] uint32_t format_version() const override { return kFormatVersion; } - void write_index(IOWriter* io_writer) override; - // Unsupported: the inline forward index is mmap-only. Throws. - void read_index(IOReader* io_reader, const IndexHeader& header, - int io_flags = 0) override; - - auto search(idx_t n, const idx_t* indptr, const term_t* indices, - const float* values, int k, - SearchParameters* search_parameters = nullptr) - -> pair_of_score_id_vectors_t override; - - // Selects the global top-k_prime blocks across `cuts` by summary score and - // scores their docs. `dense` is per-thread scratch, zero on entry and - // restored on exit; `visited` is cleared on entry. - auto single_query(std::vector& dense, - absl::flat_hash_set& visited, - const term_t* q_indices, const float* q_values, - size_t q_len, const std::vector& cuts, int k, - int k_prime, SearchParameters* search_parameters) - -> pair_of_score_id_vector_t; - - SeismicClusterParameters cluster_parameter_; - // Borrows from the base's mapped_file_, so it is declared here (destroyed - // before the mapping). Empty after a fresh build (search uses vectors_ - // then). - detail::InlineForwardIndex fwd_; - size_t num_vectors_ = 0; + [[nodiscard]] size_t code_element_size() const override; + const uint8_t* encode_values(const float* values, size_t nnz, + std::vector& scratch) const override; + const uint8_t* encode_query( + const float* values, size_t nnz, + const SearchParameters* search_parameters, + std::vector& scratch) const override; }; } // namespace nsparse diff --git a/nsparse/disk_seismic_index_base.cpp b/nsparse/disk_seismic_index_base.cpp new file mode 100644 index 0000000..4078a06 --- /dev/null +++ b/nsparse/disk_seismic_index_base.cpp @@ -0,0 +1,205 @@ +/** + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +#include "nsparse/disk_seismic_index_base.h" + +#include +#include +#include +#include +#include +#include +#include +#include + +#include "absl/container/flat_hash_set.h" +#include "nsparse/cluster/inverted_list_clusters.h" +#include "nsparse/disk_seismic_search.h" +#include "nsparse/id_selector.h" +#include "nsparse/index.h" +#include "nsparse/io/inline_forward_index_io.h" +#include "nsparse/io/seismic_invlists_writer.h" +#include "nsparse/sparse_vectors.h" +#include "nsparse/types.h" +#include "nsparse/utils/checks.h" +#include "nsparse/utils/mmap_cursor.h" +#include "nsparse/utils/mmap_file.h" + +namespace nsparse { + +DiskSeismicIndexBase::DiskSeismicIndexBase(int dim, + SeismicClusterParameters parameter) + : MmapIndex(dim), cluster_parameter_(parameter) {} + +void DiskSeismicIndexBase::add(idx_t n, const idx_t* indptr, + const term_t* indices, const float* values) { + throw_if_not_positive(n); + throw_if_any_null(indptr, indices, values); + const size_t indptr_size = n + 1; + const size_t nnz = indptr[n]; + const size_t element_size = code_element_size(); + if (vectors_ == nullptr) { + // Fresh container: start the count at 0 so a stale num_vectors_ (e.g. + // left by a prior mmap load, which has no vectors_) cannot accumulate. + num_vectors_ = 0; + vectors_ = std::make_unique(SparseVectorsConfig{ + .element_size = element_size, + .dimension = static_cast(dimension_)}); + } + std::vector scratch; + const uint8_t* codes = encode_values(values, nnz, scratch); + vectors_->add_vectors(indptr, indptr_size, indices, nnz, codes, + nnz * element_size); + num_vectors_ += n; +} + +void DiskSeismicIndexBase::build() { + clustered_inverted_lists = detail::build_inverted_lists_clusters( + get_vectors(), + {.element_size = code_element_size(), + .dimension = static_cast(get_dimension())}, + cluster_parameter_); +} + +auto DiskSeismicIndexBase::search(idx_t n, const idx_t* indptr, + const term_t* indices, const float* values, + int k, SearchParameters* search_parameters) + -> pair_of_score_id_vectors_t { + // Quit early when there is nothing to score: no vectors, no queries, or no + // forward-vector source (fwd_ empty and vectors_ null — a corrupt or + // uninitialized index). + if (num_vectors_ == 0 || n == 0 || + (fwd_.num_blocks() == 0 && vectors_ == nullptr)) { + return detail::initialize_padded_results(n, k); + } + + // The id-selector exact-match fast path is omitted (it needs an in-RAM + // SparseVectors a mapped index lacks); the selector is still honored per + // candidate doc inside the shared search core. + const detail::DiskSeismicCutBudget budget = + detail::resolve_cut_and_budget(search_parameters); + const int cut = budget.cut; + const int k_prime = budget.k_prime; + // k_prime is a block budget, not a document count: a block holds many docs, + // so k_prime < k is valid (a few blocks can still fill k, and a short-fall + // is padded like any under-budget search). Only a non-positive budget is + // rejected. + if (k_prime <= 0) { + throw std::invalid_argument( + "DiskSeismic index: k_prime (block budget) must be positive"); + } + + // Encode the whole query batch once at the stored width; a query's codes + // start at query_batch + start * element_size. + const size_t element_size = code_element_size(); + const size_t nnz = indptr[n]; + std::vector query_scratch; + const uint8_t* query_batch = + encode_query(values, nnz, search_parameters, query_scratch); + + // Rows are filled below, so start them empty rather than paying an n*k + // padding fill only to overwrite it. + std::vector> result_distances(n); + std::vector> result_labels(n); + + const detail::InlineForwardIndex* fwd = + fwd_.num_blocks() > 0 ? &fwd_ : nullptr; + const SparseVectors* vectors = fwd == nullptr ? vectors_.get() : nullptr; + const IDSelector* id_selector = search_parameters == nullptr + ? nullptr + : search_parameters->get_id_selector(); + const size_t dense_bytes = static_cast(dimension_) * element_size; + +#pragma omp parallel + { + // Per-thread scratch reused across the queries a thread handles: the + // dense lookup table, the visited-doc set, and the block-candidate / + // summary-score buffers, so no query allocates on the hot path. + std::vector dense(dense_bytes, 0); + absl::flat_hash_set visited; + visited.reserve(static_cast(std::max(k, 1)) * 4096); + std::vector candidates; + std::vector score_scratch; + +#pragma omp for schedule(dynamic, 64) + for (idx_t query_idx = 0; query_idx < n; ++query_idx) { + const idx_t start = indptr[query_idx]; + const size_t len = indptr[query_idx + 1] - start; + const term_t* query_indices = indices + start; + const uint8_t* query_codes = + query_batch + static_cast(start) * element_size; + const std::vector cuts = detail::top_cut_tokens( + query_indices, query_codes, len, cut, element_size); + auto [scores, ids] = detail::block_budget_query( + dense.data(), element_size, visited, candidates, score_scratch, + query_indices, query_codes, len, cuts, k, k_prime, + clustered_inverted_lists, fwd, vectors, id_selector); + decode_scores(scores, search_parameters); + scores.resize(k, -1.0F); + ids.resize(k, detail::INVALID_IDX); + result_distances[query_idx] = std::move(scores); + result_labels[query_idx] = std::move(ids); + } + } + return {result_distances, result_labels}; +} + +void DiskSeismicIndexBase::write_index(IOWriter* io_writer) { + write_payload_header(io_writer); + uint64_t nv = num_vectors_; + io_writer->write(&nv, sizeof(uint64_t), 1); + // Summaries only: the doc-id membership is already in the inline forward + // index below, so writing it in the posting lists too would duplicate it. + SeismicInvertedListsWriter inv_list_writer(clustered_inverted_lists, + /*summaries_only=*/true); + inv_list_writer.serialize(io_writer); + // Inline forward index, built from the same clusters + vectors. An empty + // corpus uses a correctly-typed empty SparseVectors (element_size must be a + // valid width even with zero vectors) so the section still round-trips. + SparseVectors empty_vectors( + {.element_size = code_element_size(), + .dimension = static_cast(dimension_)}); + const SparseVectors& v = vectors_ != nullptr ? *vectors_ : empty_vectors; + detail::InlineForwardIndex forward(clustered_inverted_lists, v); + forward.serialize(io_writer); +} + +void DiskSeismicIndexBase::read_index(IOReader* /*io_reader*/, + const IndexHeader& /*header*/, + int /*io_flags*/) { + // The inline forward index is borrowed from a mapping, never copied onto + // the heap, so this index has no copying read path. + throw std::runtime_error( + "DiskSeismic index is mmap-only; load with read_index(file, " + "IndexIoFlag::kUseMmap)"); +} + +void DiskSeismicIndexBase::load_mapped_payload(MmapCursor* cursor, + MmapFile&& mapped) { + // Same order write_index wrote them (past any extra header the caller + // already consumed): doc count, summaries, inline forward. + num_vectors_ = cursor->read_scalar(); + SeismicInvertedListsWriter inv_list_writer; + inv_list_writer.mmap_deserialize(cursor); + detail::InlineForwardIndex forward; + forward.mmap_deserialize(cursor); + + clustered_inverted_lists = std::move(inv_list_writer.release()); + fwd_ = std::move(forward); + // Now that the summaries and forward index are populated (still borrowing + // from `mapped`, which is alive here), let the concrete index reject a width + // mismatch before we commit. + validate_mapped_payload(); + + // mapped_file_ last: the summaries and the forward index borrow from it, and + // moving it does not move the mapping. + mapped_file_ = std::move(mapped); +} + +} // namespace nsparse diff --git a/nsparse/disk_seismic_index_base.h b/nsparse/disk_seismic_index_base.h new file mode 100644 index 0000000..8af4dc0 --- /dev/null +++ b/nsparse/disk_seismic_index_base.h @@ -0,0 +1,108 @@ +/** + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +#ifndef DISK_SEISMIC_INDEX_BASE_H +#define DISK_SEISMIC_INDEX_BASE_H + +#include +#include +#include + +#include "nsparse/cluster/inverted_list_clusters.h" +#include "nsparse/index.h" +#include "nsparse/io/inline_forward_index_io.h" +#include "nsparse/io/io.h" +#include "nsparse/mmap_index.h" +#include "nsparse/types.h" +#include "nsparse/utils/mmap_cursor.h" +#include "nsparse/utils/mmap_file.h" + +namespace nsparse { + +// Shared implementation of the two disk-resident SEISMIC indexes: the cluster +// summaries live in RAM, the per-document forward vectors live on disk in the +// block-contiguous (inline) layout and are borrowed via mmap at search time, +// and search scores the global top-k_prime blocks (a DiskSeismicSearchParameters +// sets k_prime). add / build / search / serialization / mmap loading are all +// here; the concrete indexes differ only in the stored value width and, for the +// quantized one, a leading quantization header and a score-decoding step, which +// they supply through the virtual hooks below. +// +// mmap-only: load with read_index(file, kUseMmap); the copying read throws. +class DiskSeismicIndexBase : public MmapIndex, public IndexIO { +public: + // Persisted, since a mapped index has no in-RAM vectors_ to derive it from. + size_t num_vectors() const override { return num_vectors_; } + + void add(idx_t n, const idx_t* indptr, const term_t* indices, + const float* values) override; + void build() override; + +protected: + DiskSeismicIndexBase(int dim, SeismicClusterParameters parameter); + + // --- Hooks the concrete indexes implement. --- + + // Stored value width in bytes: 4 (float) or 1/2 (quantized codes). + [[nodiscard]] virtual size_t code_element_size() const = 0; + + // Encode nnz float values to the stored width, returning a pointer to + // code_element_size()-byte-per-value data. `scratch` backs the result when + // encoding must allocate; the float index returns its input reinterpreted, + // with no copy. + virtual const uint8_t* encode_values(const float* values, size_t nnz, + std::vector& scratch) const = 0; + + // Encode a query batch the same way, honoring any per-search range override. + virtual const uint8_t* encode_query( + const float* values, size_t nnz, + const SearchParameters* search_parameters, + std::vector& scratch) const = 0; + + // Decode raw integer dot products back to float. No-op for the float index. + virtual void decode_scores( + std::vector& /*scores*/, + const SearchParameters* /*search_parameters*/) const {} + + // Extra header this index writes before the shared payload (the + // quantization header for the quantized index; nothing for the float one). + virtual void write_payload_header(IOWriter* /*io_writer*/) const {} + + // Reject a just-mapped payload whose stored width disagrees with this index. + // No-op for the float index, which stores a fixed width. Runs after the + // summaries and forward index are populated but before the mapping commits. + virtual void validate_mapped_payload() const {} + + // Reads the shared payload (doc count, summaries, inline forward) from the + // cursor, validates it, and commits the mapping. Each concrete mmap_index + // reads its own extra header first, then calls this. + void load_mapped_payload(MmapCursor* cursor, MmapFile&& mapped); + + // Borrowed from by score_summaries_transposed / the inline forward index, so + // the concrete validate_mapped_payload can inspect their widths. + std::vector clustered_inverted_lists; + detail::InlineForwardIndex fwd_; + +private: + auto search(idx_t n, const idx_t* indptr, const term_t* indices, + const float* values, int k, + SearchParameters* search_parameters = nullptr) + -> pair_of_score_id_vectors_t override; + void write_index(IOWriter* io_writer) override; + // Unsupported: the inline forward index is mmap-only. Throws. + void read_index(IOReader* io_reader, const IndexHeader& header, + int io_flags = 0) override; + + SeismicClusterParameters cluster_parameter_; + size_t num_vectors_ = 0; +}; + +} // namespace nsparse + +#endif // DISK_SEISMIC_INDEX_BASE_H diff --git a/nsparse/disk_seismic_scalar_quantized_index.cpp b/nsparse/disk_seismic_scalar_quantized_index.cpp new file mode 100644 index 0000000..7eb4201 --- /dev/null +++ b/nsparse/disk_seismic_scalar_quantized_index.cpp @@ -0,0 +1,170 @@ +/** + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +#include "nsparse/disk_seismic_scalar_quantized_index.h" + +#include +#include +#include +#include +#include +#include + +#include "nsparse/cluster/inverted_list_clusters.h" +#include "nsparse/disk_seismic_index_base.h" +#include "nsparse/index.h" +#include "nsparse/io/inline_forward_index_io.h" +#include "nsparse/types.h" +#include "nsparse/utils/checks.h" +#include "nsparse/utils/mmap_cursor.h" +#include "nsparse/utils/mmap_file.h" +#include "nsparse/utils/scalar_quantizer.h" + +namespace nsparse { +namespace { + +// Validate a quantizer type declared by a stored index and construct it. +// bytes_per_value() treats anything but QT_8bit as 16-bit, so an undefined type +// would silently pick an element width rather than be rejected. +ScalarQuantizer make_scalar_quantizer(QuantizerType type, float vmin, + float vmax) { + if (type != QuantizerType::QT_8bit && type != QuantizerType::QT_16bit) { + throw std::runtime_error( + "index file declares an unknown quantizer type"); + } + return ScalarQuantizer(type, vmin, vmax); +} + +// The stored codes are strided at the width the quantizer reports, both in the +// forward index and in the cluster summaries (score_summaries_transposed +// dispatches on the summaries' own width). A file where either disagrees would +// be read at the wrong stride, so reject it at load. +void throw_if_forward_width_mismatch(const detail::InlineForwardIndex& fwd, + const ScalarQuantizer& sq) { + if (fwd.num_blocks() > 0 && fwd.element_size() != sq.bytes_per_value()) { + throw std::runtime_error( + "index file's forward element size disagrees with its quantizer " + "type"); + } +} + +void throw_if_summary_width_mismatch( + const std::vector& clusters, + const ScalarQuantizer& sq) { + for (const InvertedListClusters& list : clusters) { + if (list.cluster_size() > 0 && + list.element_size() != sq.bytes_per_value()) { + throw std::runtime_error( + "index file's summary element size disagrees with its " + "quantizer type"); + } + } +} + +} // namespace + +DiskSeismicScalarQuantizedIndex::DiskSeismicScalarQuantizedIndex(int dim) + : DiskSeismicIndexBase(dim, detail::kDefaultSeismicClusterParams) {} + +DiskSeismicScalarQuantizedIndex::DiskSeismicScalarQuantizedIndex( + QuantizerType quantizer_type, float vmin, float vmax, + SeismicClusterParameters parameter, int dim) + : DiskSeismicIndexBase(dim, parameter), sq_(quantizer_type, vmin, vmax) {} + +void DiskSeismicScalarQuantizedIndex::read_csr(const char* file_path, + Residency residency) { + if (residency == Residency::kMmap) { + throw std::invalid_argument( + "mmap residency is not available for a quantized index: a mapped " + "CSR is borrowed as float, and this index searches over codes"); + } + MmapIndex::read_csr(file_path, residency); +} + +size_t DiskSeismicScalarQuantizedIndex::code_element_size() const { + return sq_.bytes_per_value(); +} + +const uint8_t* DiskSeismicScalarQuantizedIndex::encode_values( + const float* values, size_t nnz, std::vector& scratch) const { + scratch.resize(nnz * sq_.bytes_per_value()); + sq_.encode(values, scratch.data(), nnz); + return scratch.data(); +} + +const uint8_t* DiskSeismicScalarQuantizedIndex::encode_query( + const float* values, size_t nnz, + const SearchParameters* search_parameters, + std::vector& scratch) const { + const ScalarQuantizer query_sq = query_quantizer(search_parameters); + scratch.resize(nnz * sq_.bytes_per_value()); + query_sq.encode(values, scratch.data(), nnz); + return scratch.data(); +} + +void DiskSeismicScalarQuantizedIndex::decode_scores( + std::vector& scores, + const SearchParameters* search_parameters) const { + // Decode the integer dot products back to float. Called before the -1.0 + // padding is added, so the sentinel is never scaled. + const ScalarQuantizer query_sq = query_quantizer(search_parameters); + for (float& score : scores) { + score = sq_.decode_dot_product(score, query_sq); + } +} + +// The type is always the index's, since the codes are compared against stored +// ones; only the range can be overridden per search. +ScalarQuantizer DiskSeismicScalarQuantizedIndex::query_quantizer( + const SearchParameters* search_parameters) const { + const auto* sq_params = + dynamic_cast(search_parameters); + if (sq_params == nullptr) { + return sq_; + } + return ScalarQuantizer(sq_.get_quantizer_type(), sq_params->vmin, + sq_params->vmax); +} + +void DiskSeismicScalarQuantizedIndex::write_payload_header( + IOWriter* io_writer) const { + auto sq_type = sq_.get_quantizer_type(); + io_writer->write(&sq_type, sizeof(QuantizerType), 1); + auto vmin = sq_.get_min(); + io_writer->write(&vmin, sizeof(float), 1); + auto vmax = sq_.get_max(); + io_writer->write(&vmax, sizeof(float), 1); +} + +void DiskSeismicScalarQuantizedIndex::validate_mapped_payload() const { + throw_if_forward_width_mismatch(fwd_, sq_); + throw_if_summary_width_mismatch(clustered_inverted_lists, sq_); +} + +DiskSeismicScalarQuantizedIndex* DiskSeismicScalarQuantizedIndex::mmap_index( + const IndexHeader& header, const char* index_file, size_t pos) { + throw_if_null(index_file, "index_file must not be null"); + auto index = + std::make_unique(header.dimension); + + MmapFile mmap_file(std::string{index_file}); + MmapCursor cursor(mmap_file.data(), mmap_file.size()); + cursor.skip(pos); + + // The quantization header opens the payload; the shared loader reads the + // rest and calls validate_mapped_payload() against this quantizer. + const auto sq_type = cursor.read_scalar(); + const auto vmin = cursor.read_scalar(); + const auto vmax = cursor.read_scalar(); + index->sq_ = make_scalar_quantizer(sq_type, vmin, vmax); + + index->load_mapped_payload(&cursor, std::move(mmap_file)); + return index.release(); +} +} // namespace nsparse diff --git a/nsparse/disk_seismic_scalar_quantized_index.h b/nsparse/disk_seismic_scalar_quantized_index.h new file mode 100644 index 0000000..fd344da --- /dev/null +++ b/nsparse/disk_seismic_scalar_quantized_index.h @@ -0,0 +1,104 @@ +/** + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +#ifndef DISK_SEISMIC_SCALAR_QUANTIZED_INDEX_H +#define DISK_SEISMIC_SCALAR_QUANTIZED_INDEX_H +#include +#include +#include +#include + +#include "nsparse/disk_seismic_index.h" // DiskSeismicSearchParameters +#include "nsparse/disk_seismic_index_base.h" +#include "nsparse/io/io.h" +#include "nsparse/mmap_index.h" // Residency +#include "nsparse/types.h" +#include "nsparse/utils/scalar_quantizer.h" + +namespace nsparse { + +// Overrides the quantizer range a query is encoded with (the type is always the +// index's); everything else is DiskSeismicSearchParameters (cut + k_prime). +struct DiskSeismicSQSearchParameters : public DiskSeismicSearchParameters { + float vmin; + float vmax; + DiskSeismicSQSearchParameters(float vmin, float vmax, int cut, int k_prime) + : DiskSeismicSearchParameters(cut, k_prime), vmax(vmax), vmin(vmin) {} +}; + +// A DiskSeismicIndex over scalar-quantized codes: the forward vectors and the +// cluster summaries are stored as 8- or 16-bit codes instead of float. The +// disk-resident search and serialization are shared with DiskSeismicIndex via +// DiskSeismicIndexBase; this type supplies the code width, the query +// quantization + score decoding, and a leading quantization header. Pass a +// DiskSeismicSQSearchParameters to override the query range, else the index's +// build-time range is reused. +// +// mmap-only: load with read_index(file, kUseMmap); the copying read throws. +class DiskSeismicScalarQuantizedIndex : public DiskSeismicIndexBase { +public: + static constexpr std::array name = {'D', 'S', 'S', 'Q'}; + // Bump whenever write_index's payload layout changes. + static constexpr uint32_t kFormatVersion = 1; + + explicit DiskSeismicScalarQuantizedIndex(int dim); + DiskSeismicScalarQuantizedIndex(QuantizerType quantizer_type, float vmin, + float vmax, + SeismicClusterParameters parameter, + int dim); + ~DiskSeismicScalarQuantizedIndex() override = default; + std::array id() const override { return name; } + + DiskSeismicScalarQuantizedIndex(const DiskSeismicScalarQuantizedIndex&) = + delete; + DiskSeismicScalarQuantizedIndex& operator=( + const DiskSeismicScalarQuantizedIndex&) = delete; + + const ScalarQuantizer& get_scalar_quantizer() const { return sq_; } + + // Borrows a serialized index from a file mapping. `pos` is where the + // payload begins, past the header read_header consumed. + static DiskSeismicScalarQuantizedIndex* mmap_index(const IndexHeader& header, + const char* index_file, + size_t pos); + + // Only the copying residency. A mapped CSR is borrowed at the width it was + // written in, which is float, whereas this index searches over codes: the + // values have to pass through add() to be quantized. + void read_csr(const char* file_path, + Residency residency = Residency::kInMemory) override; + +private: + [[nodiscard]] uint32_t format_version() const override { + return kFormatVersion; + } + [[nodiscard]] size_t code_element_size() const override; + const uint8_t* encode_values(const float* values, size_t nnz, + std::vector& scratch) const override; + const uint8_t* encode_query( + const float* values, size_t nnz, + const SearchParameters* search_parameters, + std::vector& scratch) const override; + void decode_scores(std::vector& scores, + const SearchParameters* search_parameters) const override; + // The quantization parameters that open this index's payload -- distinct + // from the IndexHeader the file itself starts with. + void write_payload_header(IOWriter* io_writer) const override; + void validate_mapped_payload() const override; + + // The quantizer a query is encoded with: DiskSeismicSQSearchParameters + // overrides the range the index was built with, anything else reuses it. + ScalarQuantizer query_quantizer( + const SearchParameters* search_parameters) const; + + ScalarQuantizer sq_; +}; +} // namespace nsparse + +#endif // DISK_SEISMIC_SCALAR_QUANTIZED_INDEX_H diff --git a/nsparse/disk_seismic_search.cpp b/nsparse/disk_seismic_search.cpp new file mode 100644 index 0000000..aeffad9 --- /dev/null +++ b/nsparse/disk_seismic_search.cpp @@ -0,0 +1,207 @@ +/** + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +#include "nsparse/disk_seismic_search.h" + +#include +#include +#include +#include + +#include "absl/container/flat_hash_set.h" +#include "nsparse/cluster/inverted_list_clusters.h" +#include "nsparse/id_selector.h" +#include "nsparse/index.h" +#include "nsparse/io/inline_forward_index_io.h" +#include "nsparse/seismic_index.h" // SeismicSearchParameters +#include "nsparse/sparse_vectors.h" +#include "nsparse/types.h" +#include "nsparse/utils/distance_simd.h" +#include "nsparse/utils/ranker.h" +#include "nsparse/utils/vector_process.h" + +namespace nsparse::detail { +namespace { + +// Score one doc's vector against the dense query and push it, honoring +// visited-dedup and the id selector. `dense` and `vals` are element_size bytes +// per component; the width selects the kernel (4 = float, 2 = uint16, else +// uint8). The score is the raw dot product; a quantized caller decodes later. +inline void score_doc(idx_t doc_id, const term_t* comps, const uint8_t* vals, + size_t nnz, const uint8_t* dense, size_t element_size, + const IDSelector* id_selector, TopKHolder& heap, + absl::flat_hash_set& visited) { + if (!visited.insert(doc_id).second) { + return; + } + if (id_selector != nullptr && !id_selector->is_member(doc_id)) { + return; + } + float score = 0.0F; + if (element_size == U32) { + score = dot_product_float_dense( + comps, reinterpret_cast(vals), nnz, + reinterpret_cast(dense)); + } else if (element_size == U16) { + score = dot_product_uint16_dense( + comps, reinterpret_cast(vals), nnz, + reinterpret_cast(dense)); + } else { + score = dot_product_uint8_dense(comps, vals, nnz, dense); + } + heap.add(score, doc_id); +} + +// Score every document of one block against the dense query, dedup via +// `visited` and honor the id selector. Vectors come from the inline forward +// index (fwd) when loaded, else from the in-RAM CSR of a fresh build. +void score_block(const InlineForwardIndex* fwd, const SparseVectors* vectors, + const std::vector& clusters, uint32_t pl, + uint32_t cid, const uint8_t* dense, size_t element_size, + const IDSelector* id_selector, TopKHolder& heap, + absl::flat_hash_set& visited) { + if (fwd != nullptr) { + const BlockView bv = fwd->block(pl, cid); + if (bv.absent()) { + return; + } + for (uint32_t i = 0; i < bv.n_docs; ++i) { + score_doc(static_cast(bv.doc_ids[i]), bv.doc_comps(i), + bv.doc_vals(i, element_size), bv.nnz(i), dense, + element_size, id_selector, heap, visited); + } + } else if (vectors != nullptr) { + const idx_t* const indptr = vectors->indptr_data(); + const term_t* const indices = vectors->indices_data(); + const uint8_t* const values = vectors->values_data(); + for (const idx_t doc_id : clusters[pl].get_docs(cid)) { + const idx_t start = indptr[doc_id]; + const size_t len = indptr[doc_id + 1] - start; + score_doc(doc_id, indices + start, + values + static_cast(start) * element_size, len, + dense, element_size, id_selector, heap, visited); + } + } +} + +} // namespace + +DiskSeismicCutBudget resolve_cut_and_budget( + const SearchParameters* search_parameters) { + // A DiskSeismicSearchParameters (including a quantized subclass) carries + // k_prime; a plain SeismicSearchParameters (or null) uses the default + // budget. + const DiskSeismicSearchParameters defaults; + int cut = defaults.cut; + int k_prime = defaults.k_prime; + if (const auto* disk_parameters = + dynamic_cast( + search_parameters)) { + cut = disk_parameters->cut; + k_prime = disk_parameters->k_prime; + } else if (const auto* seismic_parameters = + dynamic_cast( + search_parameters)) { + cut = seismic_parameters->cut; + } + return {cut, k_prime}; +} + +std::vector top_cut_tokens(const term_t* indices, const uint8_t* codes, + size_t len, int cut, size_t element_size) { + if (element_size == U32) { + return top_k_tokens( + indices, reinterpret_cast(codes), len, cut); + } + if (element_size == U16) { + return top_k_tokens( + indices, reinterpret_cast(codes), len, cut); + } + return top_k_tokens(indices, codes, len, cut); +} + +pair_of_score_id_vectors_t initialize_padded_results(idx_t n, int k) { + return {std::vector>(n, std::vector(k, -1.0F)), + std::vector>( + n, std::vector(k, INVALID_IDX))}; +} + +pair_of_score_id_vector_t block_budget_query( + uint8_t* dense, size_t element_size, absl::flat_hash_set& visited, + std::vector& candidates, std::vector& score_scratch, + const term_t* query_indices, const uint8_t* query_values, size_t query_len, + const std::vector& cuts, int k, int k_prime, + const std::vector& clusters, + const InlineForwardIndex* fwd, const SparseVectors* vectors, + const IDSelector* id_selector) { + // Scatter the query into the dense lookup table: element_size contiguous + // bytes per non-zero dim. + for (size_t i = 0; i < query_len; ++i) { + std::copy_n( + query_values + i * element_size, element_size, + dense + static_cast(query_indices[i]) * element_size); + } + visited.clear(); + + // Collect every candidate block across the cut lists with its summary + // score. score_summaries_transposed dots the query with each cluster + // summary (stored at the same width), so scores are comparable across + // posting lists for the global ranking. + candidates.clear(); + for (const term_t term : cuts) { + if (term >= clusters.size()) [[unlikely]] { + continue; + } + const InvertedListClusters& cluster_invlist = clusters[term]; + const size_t n_clusters = cluster_invlist.cluster_size(); + if (n_clusters == 0) { + continue; + } + cluster_invlist.score_summaries_transposed(query_indices, query_values, + query_len, score_scratch); + for (size_t cid = 0; cid < n_clusters; ++cid) { + candidates.push_back( + {score_scratch[cid], term, static_cast(cid)}); + } + } + + // Select the global top-k_prime blocks by summary score. The order among + // them does not affect the result (dedup + exact scoring are + // order-independent). + const size_t budget = + std::min(static_cast(k_prime), candidates.size()); + if (budget < candidates.size()) { + std::nth_element(candidates.begin(), candidates.begin() + budget, + candidates.end(), + [](const BlockCandidate& a, const BlockCandidate& b) { + return a.score > b.score; + }); + candidates.resize(budget); + } + + // Score the selected blocks; visited dedups docs shared across them. + TopKHolder holder(k); + for (const BlockCandidate& candidate : candidates) { + score_block(fwd, vectors, clusters, candidate.pl, candidate.cid, dense, + element_size, id_selector, holder, visited); + } + + // Restore only the query's own positions to zero (mirrors the scatter at + // entry): element_size bytes per touched dim. The scattered writes are + // bounded by query_len and negligible next to the block scoring. + for (size_t i = 0; i < query_len; ++i) { + std::fill_n( + dense + static_cast(query_indices[i]) * element_size, + element_size, uint8_t{0}); + } + + return holder.top_k_items_descending(); +} + +} // namespace nsparse::detail diff --git a/nsparse/disk_seismic_search.h b/nsparse/disk_seismic_search.h new file mode 100644 index 0000000..343b525 --- /dev/null +++ b/nsparse/disk_seismic_search.h @@ -0,0 +1,79 @@ +/** + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +#ifndef DISK_SEISMIC_SEARCH_H +#define DISK_SEISMIC_SEARCH_H + +#include +#include +#include + +#include "absl/container/flat_hash_set.h" +#include "nsparse/cluster/inverted_list_clusters.h" +#include "nsparse/disk_seismic_index.h" // DiskSeismicSearchParameters +#include "nsparse/id_selector.h" +#include "nsparse/index.h" // SearchParameters, pair_of_score_id_vector(s)_t +#include "nsparse/io/inline_forward_index_io.h" +#include "nsparse/sparse_vectors.h" +#include "nsparse/types.h" + +// Search machinery shared by the two disk-resident SEISMIC indexes, +// DiskSeismicIndex (float) and DiskSeismicScalarQuantizedIndex (8-/16-bit +// codes). They differ only in the value width, so the block-budget search is +// one element_size-parameterized routine here rather than a copy in each. +namespace nsparse::detail { + +// A candidate block: its summary score and its (posting list, cluster) address. +struct BlockCandidate { + float score; + uint32_t pl; + uint32_t cid; +}; + +// (cut, k_prime) resolved from search parameters, defaulting to +// DiskSeismicSearchParameters. Each index still validates k_prime with its own +// error message. +struct DiskSeismicCutBudget { + int cut; + int k_prime; +}; +DiskSeismicCutBudget resolve_cut_and_budget( + const SearchParameters* search_parameters); + +// The `cut` highest-weighted query terms, read from `codes` at `element_size` +// bytes per value (4 = float, 2 = uint16, 1 = uint8). +std::vector top_cut_tokens(const term_t* indices, const uint8_t* codes, + size_t len, int cut, size_t element_size); + +// The all-padding result grid a search returns when there is nothing to score. +// Distances -1.0, labels INVALID_IDX. +pair_of_score_id_vectors_t initialize_padded_results(idx_t n, int k); + +// Block-budget query core. The query must already be encoded at `element_size` +// bytes per component in `query_values` -- raw float (width 4) for the plain +// index, quantized codes (width 1 or 2) for the quantized one; the width also +// selects the dot-product kernel. `dense` is a dimension*element_size byte +// buffer, all-zero on entry and restored to all-zero on exit; `visited`, +// `candidates`, and `score_scratch` are per-thread scratch reused across +// queries (cleared/overwritten here). Blocks come from `fwd` (the mapping) when +// it is non-null, else from the in-RAM `vectors` of a fresh build. Returns the +// raw (undecoded) top-k scores + ids, unpadded: the caller decodes if it must +// and pads to k. +pair_of_score_id_vector_t block_budget_query( + uint8_t* dense, size_t element_size, absl::flat_hash_set& visited, + std::vector& candidates, std::vector& score_scratch, + const term_t* query_indices, const uint8_t* query_values, size_t query_len, + const std::vector& cuts, int k, int k_prime, + const std::vector& clusters, + const InlineForwardIndex* fwd, const SparseVectors* vectors, + const IDSelector* id_selector); + +} // namespace nsparse::detail + +#endif // DISK_SEISMIC_SEARCH_H diff --git a/nsparse/index_factory.cpp b/nsparse/index_factory.cpp index 7fcf35d..2beeef6 100644 --- a/nsparse/index_factory.cpp +++ b/nsparse/index_factory.cpp @@ -17,6 +17,7 @@ #include "nsparse/brutal_index.h" #include "nsparse/disk_seismic_index.h" +#include "nsparse/disk_seismic_scalar_quantized_index.h" #include "nsparse/id_map_index.h" #include "nsparse/index.h" #include "nsparse/inverted_index.h" @@ -50,6 +51,41 @@ std::string trim(const std::string& str) { return str.substr(first, last - first + 1); } +// The cluster knobs shared by every seismic-family index. `get_param` is the +// (key, default) -> value lookup over the parsed description. +template +SeismicClusterParameters parse_cluster_params(const GetParam& get_param) { + return {.lambda = std::stoi(get_param("lambda", "10")), + .beta = std::stoi(get_param("beta", "5")), + .alpha = std::stof(get_param("alpha", "0.5")), + .seed = std::stoi(get_param("seed", std::to_string(kRandomSeed)))}; +} + +struct QuantizerConfig { + QuantizerType type; + float vmin; + float vmax; +}; + +// The quantizer knobs shared by seismic_sq and disk_seismic_sq. An unrecognized +// quantizer= is rejected rather than silently falling back to a width. +template +QuantizerConfig parse_quantizer_config(const GetParam& get_param) { + const std::string quantizer_str = get_param("quantizer", "8bit"); + QuantizerType type = QuantizerType::QT_8bit; + if (quantizer_str == "8bit") { + type = QuantizerType::QT_8bit; + } else if (quantizer_str == "16bit") { + type = QuantizerType::QT_16bit; + } else { + throw std::invalid_argument("unknown quantizer '" + quantizer_str + + "'; expected '8bit' or '16bit'"); + } + return {.type = type, + .vmin = std::stof(get_param("vmin", "0.0")), + .vmax = std::stof(get_param("vmax", "1.0"))}; +} + } // namespace // description is segmented by , @@ -102,44 +138,23 @@ Index* index_factory(int dimension, const char* description) { } if (index_type == "seismic") { - int lambda = std::stoi(get_param("lambda", "10")); - int beta = std::stoi(get_param("beta", "5")); - float alpha = std::stof(get_param("alpha", "0.5")); - int seed = std::stoi(get_param("seed", std::to_string(kRandomSeed))); - return new SeismicIndex(dimension, {.lambda = lambda = lambda, - .beta = beta = beta, - .alpha = alpha = alpha, - .seed = seed}); + return new SeismicIndex(dimension, parse_cluster_params(get_param)); } if (index_type == "disk_seismic") { - int lambda = std::stoi(get_param("lambda", "10")); - int beta = std::stoi(get_param("beta", "5")); - float alpha = std::stof(get_param("alpha", "0.5")); - int seed = std::stoi(get_param("seed", std::to_string(kRandomSeed))); - return new DiskSeismicIndex( - dimension, - {.lambda = lambda, .beta = beta, .alpha = alpha, .seed = seed}); + return new DiskSeismicIndex(dimension, parse_cluster_params(get_param)); + } + + if (index_type == "disk_seismic_sq") { + const QuantizerConfig q = parse_quantizer_config(get_param); + return new DiskSeismicScalarQuantizedIndex( + q.type, q.vmin, q.vmax, parse_cluster_params(get_param), dimension); } if (index_type == "seismic_sq") { - std::string quantizer_str = get_param("quantizer", "8bit"); - QuantizerType quantizer_type = QuantizerType::QT_8bit; - if (quantizer_str == "16bit") { - quantizer_type = QuantizerType::QT_16bit; - } - float vmin = std::stof(get_param("vmin", "0.0")); - float vmax = std::stof(get_param("vmax", "1.0")); - int lambda = std::stoi(get_param("lambda", "10")); - int beta = std::stoi(get_param("beta", "5")); - float alpha = std::stof(get_param("alpha", "0.5")); - int seed = std::stoi(get_param("seed", std::to_string(kRandomSeed))); - return new SeismicScalarQuantizedIndex(quantizer_type, vmin, vmax, - {.lambda = lambda = lambda, - .beta = beta = beta, - .alpha = alpha = alpha, - .seed = seed}, - dimension); + const QuantizerConfig q = parse_quantizer_config(get_param); + return new SeismicScalarQuantizedIndex( + q.type, q.vmin, q.vmax, parse_cluster_params(get_param), dimension); } if (index_type == "idmap") { diff --git a/nsparse/io/index_io.cpp b/nsparse/io/index_io.cpp index 1afc9a3..d1ad869 100644 --- a/nsparse/io/index_io.cpp +++ b/nsparse/io/index_io.cpp @@ -17,6 +17,7 @@ #include "nsparse/brutal_index.h" #include "nsparse/disk_seismic_index.h" +#include "nsparse/disk_seismic_scalar_quantized_index.h" #include "nsparse/id_map_index.h" #include "nsparse/inverted_index.h" #include "nsparse/io/file_io.h" @@ -34,6 +35,7 @@ constexpr uint32_t SESQ = fourcc(SeismicScalarQuantizedIndex::name); constexpr uint32_t IDMP = fourcc(IDMapIndex::name); constexpr uint32_t INVT = fourcc(InvertedIndex::name); constexpr uint32_t DSEI = fourcc(DiskSeismicIndex::name); +constexpr uint32_t DSSQ = fourcc(DiskSeismicScalarQuantizedIndex::name); // Closes a stream once, on whichever path leaves the scope. // @@ -90,6 +92,9 @@ Index* mmap_index_payload(const IndexHeader& header, const char* file_name, return InvertedIndex::mmap_index(header, file_name, pos); case DSEI: return DiskSeismicIndex::mmap_index(header, file_name, pos); + case DSSQ: + return DiskSeismicScalarQuantizedIndex::mmap_index(header, + file_name, pos); default: return nullptr; } @@ -138,6 +143,8 @@ Index* make_index(const IndexHeader& header) { return new SeismicScalarQuantizedIndex(header.dimension); case DSEI: return new DiskSeismicIndex(header.dimension); + case DSSQ: + return new DiskSeismicScalarQuantizedIndex(header.dimension); case IDMP: return new IDMapIndex(); case INVT: diff --git a/nsparse/python/swignsparse.swig b/nsparse/python/swignsparse.swig index b5bd011..109a8d6 100644 --- a/nsparse/python/swignsparse.swig +++ b/nsparse/python/swignsparse.swig @@ -25,7 +25,9 @@ #include "nsparse/seismic_common.h" #include "nsparse/seismic_index.h" #include "nsparse/seismic_scalar_quantized_index.h" +#include "nsparse/disk_seismic_index_base.h" #include "nsparse/disk_seismic_index.h" +#include "nsparse/disk_seismic_scalar_quantized_index.h" #include "nsparse/id_map_index.h" #include "nsparse/utils/scalar_quantizer.h" #include "nsparse/cluster/random_kmeans.h" @@ -140,6 +142,7 @@ import_array(); %ignore nsparse::SeismicScalarQuantizedIndex::mmap_index; %ignore nsparse::InvertedIndex::mmap_index; %ignore nsparse::DiskSeismicIndex::mmap_index; +%ignore nsparse::DiskSeismicScalarQuantizedIndex::mmap_index; %ignore nsparse::InvertedListClusters::serialize; %ignore nsparse::InvertedListClusters::deserialize; @@ -190,7 +193,12 @@ import_array(); %include "nsparse/seismic_common.h" %include "nsparse/seismic_index.h" %include "nsparse/seismic_scalar_quantized_index.h" +// The disk-resident indexes derive from DiskSeismicIndexBase (itself a +// MmapIndex), so SWIG needs the base first or the Python classes lose their +// inherited search / num_vectors / add. +%include "nsparse/disk_seismic_index_base.h" %include "nsparse/disk_seismic_index.h" +%include "nsparse/disk_seismic_scalar_quantized_index.h" %include "nsparse/inverted_index.h" %include "nsparse/id_map_index.h" %include "nsparse/cluster/random_kmeans.h" @@ -198,10 +206,10 @@ import_array(); // index_factory returns a new object - Python takes ownership // Use %factory to enable proper downcasting based on runtime type %newobject nsparse::index_factory; -%factory(nsparse::Index* nsparse::index_factory, nsparse::BrutalIndex, nsparse::SeismicIndex, nsparse::SeismicScalarQuantizedIndex, nsparse::IDMapIndex, nsparse::InvertedIndex); +%factory(nsparse::Index* nsparse::index_factory, nsparse::BrutalIndex, nsparse::SeismicIndex, nsparse::SeismicScalarQuantizedIndex, nsparse::DiskSeismicIndex, nsparse::DiskSeismicScalarQuantizedIndex, nsparse::IDMapIndex, nsparse::InvertedIndex); %include "nsparse/index_factory.h" // read_index returns a new object - Python takes ownership %newobject nsparse::read_index; -%factory(nsparse::Index* nsparse::read_index, nsparse::BrutalIndex, nsparse::SeismicIndex, nsparse::SeismicScalarQuantizedIndex, nsparse::IDMapIndex, nsparse::InvertedIndex); +%factory(nsparse::Index* nsparse::read_index, nsparse::BrutalIndex, nsparse::SeismicIndex, nsparse::SeismicScalarQuantizedIndex, nsparse::DiskSeismicIndex, nsparse::DiskSeismicScalarQuantizedIndex, nsparse::IDMapIndex, nsparse::InvertedIndex); %include "nsparse/io/index_io.h" diff --git a/nsparse/seismic_scalar_quantized_index.cpp b/nsparse/seismic_scalar_quantized_index.cpp index f7313c5..b15f9a5 100644 --- a/nsparse/seismic_scalar_quantized_index.cpp +++ b/nsparse/seismic_scalar_quantized_index.cpp @@ -357,7 +357,7 @@ auto SeismicScalarQuantizedIndex::single_query( } void SeismicScalarQuantizedIndex::write_index(IOWriter* io_writer) { - write_quantizer_header(io_writer); + write_quantization_header(io_writer); // write vectors if (vectors_ == nullptr) { empty_sparse_vectors.serialize(io_writer); @@ -371,7 +371,7 @@ void SeismicScalarQuantizedIndex::write_index(IOWriter* io_writer) { void SeismicScalarQuantizedIndex::read_index(IOReader* io_reader, const IndexHeader& header, int io_flags) { - read_quantizer_header(io_reader); + read_quantization_header(io_reader); SparseVectors tmp_vectors; tmp_vectors.deserialize(io_reader); throw_if_element_size_mismatch(tmp_vectors, sq_); @@ -421,7 +421,8 @@ SeismicScalarQuantizedIndex* SeismicScalarQuantizedIndex::mmap_index( return index.release(); } -void SeismicScalarQuantizedIndex::write_quantizer_header(IOWriter* io_writer) { +void SeismicScalarQuantizedIndex::write_quantization_header( + IOWriter* io_writer) { auto sq_type = sq_.get_quantizer_type(); io_writer->write(&sq_type, sizeof(QuantizerType), 1); auto vmin = sq_.get_min(); @@ -430,7 +431,7 @@ void SeismicScalarQuantizedIndex::write_quantizer_header(IOWriter* io_writer) { io_writer->write(&vmax, sizeof(float), 1); } -void SeismicScalarQuantizedIndex::read_quantizer_header(IOReader* io_reader) { +void SeismicScalarQuantizedIndex::read_quantization_header(IOReader* io_reader) { QuantizerType sq_type = QuantizerType::QT_8bit; float vmin = 0.0F; float vmax = 1.0F; diff --git a/nsparse/seismic_scalar_quantized_index.h b/nsparse/seismic_scalar_quantized_index.h index 6415541..c58b0ff 100644 --- a/nsparse/seismic_scalar_quantized_index.h +++ b/nsparse/seismic_scalar_quantized_index.h @@ -79,8 +79,8 @@ class SeismicScalarQuantizedIndex : public MmapIndex, public IndexIO { int io_flags = 0) override; // The quantizer parameters that open this index's payload -- distinct from // the IndexHeader the file itself starts with. - void write_quantizer_header(IOWriter* io_writer); - void read_quantizer_header(IOReader* io_reader); + void write_quantization_header(IOWriter* io_writer); + void read_quantization_header(IOReader* io_reader); // Null `search_parameters` searches with the defaults, as it does for // SeismicIndex and as the base signature's default argument implies. This diff --git a/python_tests/disk_seismic_contract.py b/python_tests/disk_seismic_contract.py new file mode 100644 index 0000000..c617c9b --- /dev/null +++ b/python_tests/disk_seismic_contract.py @@ -0,0 +1,208 @@ +# Copyright OpenSearch Contributors +# SPDX-License-Identifier: Apache-2.0 +# +# The OpenSearch Contributors require contributions made to +# this file be licensed under the Apache-2.0 license or a +# compatible open source license. + +"""Shared black-box contract for the disk-resident SEISMIC indexes. + +DiskSeismicIndex (float) and DiskSeismicScalarQuantizedIndex (codes) present the +same public contract -- mmap-only reads, the top-k' block budget, bit-identical +fresh-build vs mmap-reload results -- so those tests live here once. A suite +subclasses `DiskSeismicContract`, sets SPEC / SEEDED / RECALL_FLOOR, and +implements `params()`; the quantized suite adds its own quantization tests. +""" + +import numpy as np +import pytest + +import nsparse +from oracle import recall_at_k +from support import ( + K, + PAD_DIST, + PAD_LABEL, + make_index, + roundtrip, + search, + search_each, + slice_corpus, +) + + +class DiskSeismicContract: + # Subclasses set these. + SPEC: str + SEEDED: str # SPEC + a fixed seed, for the determinism-sensitive tests + RECALL_FLOOR: float + + # CUT covers every query term (QUERY_NNZ=8); K_PRIME is generous so the + # block budget is not the recall bottleneck. + CUT = 8 + K_PRIME = 200 + + def params(self, cut=None, k_prime=None): + """The index's SearchParameters. Subclass provides the concrete type; + None means "use the default", so k_prime=0 is passed through (not + coalesced) for the rejection test.""" + raise NotImplementedError + + @pytest.mark.parametrize("residency", ["memory", "mmap"]) + def test_happy_case(self, residency, corpus, queries, oracle, tmp_path): + """factory -> ingest -> build -> query -> accuracy, in both residencies.""" + index = make_index(self.SPEC, corpus) + assert index.num_vectors() == corpus.n + assert index.get_dimension() == corpus.dim + + if residency == "mmap": + # disk-resident indexes are mmap-only, so the copying flag (0) would + # raise; the mmap flag is required, and the count must survive. + index = roundtrip(index, tmp_path / "index.idx", nsparse.kUseMmap) + assert index.num_vectors() == corpus.n + + dists, labels = search(index, queries, params=self.params()) + assert labels.shape == (queries.n, K) + assert dists.shape == (queries.n, K) + assert (labels[:, 0] >= 0).all(), "every query must return at least one hit" + + want_labels, _ = oracle + assert recall_at_k(labels, want_labels) >= self.RECALL_FLOOR + + def test_fresh_build_matches_mmap_reload(self, corpus, queries, tmp_path): + """The in-RAM build (vectors_) and the mmap reload (inline forward index) + are two different code paths that must return identical results.""" + index = make_index(self.SPEC, corpus) + p = self.params() + before_d, before_l = search(index, queries, params=p) + reloaded = roundtrip(index, tmp_path / "index.idx", nsparse.kUseMmap) + after_d, after_l = search(reloaded, queries, params=p) + np.testing.assert_array_equal(after_l, before_l) + np.testing.assert_allclose(after_d, before_d, rtol=1e-6, atol=1e-6) + + def test_copying_read_throws(self, corpus, tmp_path): + """mmap-only: reloading without kUseMmap must raise.""" + index = make_index(self.SPEC, corpus) + path = tmp_path / "index.idx" + nsparse.write_index(index, str(path)) + with pytest.raises(RuntimeError, match="mmap-only"): + nsparse.read_index(str(path), 0) + + @pytest.mark.parametrize("bad", [0, -1]) + def test_rejects_non_positive_block_budget(self, corpus, queries, bad): + """A non-positive k' (block budget) is rejected at search time.""" + index = make_index(self.SPEC, corpus) + with pytest.raises(ValueError, match="must be positive"): + search(index, queries, params=self.params(k_prime=bad)) + + def test_block_budget_is_monotone(self, corpus, queries, oracle, tmp_path): + """A larger block budget scores a superset of blocks, so recall never + regresses and a big enough budget beats the smallest one. Runs on a + seeded, mmap-reloaded index (deterministic on-disk block-read path).""" + seeded = roundtrip( + make_index(self.SEEDED, corpus), tmp_path / "seeded.idx", nsparse.kUseMmap + ) + want_labels, _ = oracle + recalls = [ + recall_at_k( + search(seeded, queries, params=self.params(k_prime=kp))[1], want_labels + ) + for kp in [1, 4, 16, 64, 256] + ] + for lo, hi in zip(recalls, recalls[1:]): + assert hi >= lo - 1e-9, f"recall regressed across k': {recalls}" + assert recalls[-1] > recalls[0], "a larger budget must eventually help" + + def test_block_budget_saturates(self, corpus, queries, tmp_path): + """Past the candidate-block count, more budget changes nothing.""" + seeded = roundtrip( + make_index(self.SEEDED, corpus), tmp_path / "seeded.idx", nsparse.kUseMmap + ) + big_d, big_l = search(seeded, queries, params=self.params(k_prime=10**6)) + bigger_d, bigger_l = search(seeded, queries, params=self.params(k_prime=2 * 10**6)) + np.testing.assert_array_equal(bigger_l, big_l) + np.testing.assert_array_equal(bigger_d, big_d) + + def test_empty_index_roundtrip(self, queries, tmp_path): + """An un-built, empty index writes, mmap-reloads, and returns all padding.""" + empty = nsparse.index_factory(queries.dim, self.SPEC) + path = tmp_path / "empty.idx" + nsparse.write_index(empty, str(path)) + mapped = nsparse.read_index(str(path), nsparse.kUseMmap) + assert mapped.num_vectors() == 0 + dists, labels = search(mapped, queries, params=self.params()) + assert (labels == PAD_LABEL).all() + assert (dists == PAD_DIST).all() + + def test_search_before_build(self, corpus, queries): + """Searching an added-but-unbuilt index yields empty results, not an + error. Pinned deliberately: it is a silent-empty footgun.""" + index = nsparse.index_factory(corpus.dim, self.SPEC) + index.add(corpus.n, corpus.indptr, corpus.indices, corpus.values) + _, labels = search(index, queries, params=self.params()) + assert (labels == PAD_LABEL).all() + + @pytest.mark.parametrize("residency", ["memory", "mmap"]) + def test_filtered_search(self, residency, corpus, queries, oracle, tmp_path): + """An id selector larger than k filters results to its members, on both + the in-RAM (vectors_) and mmap (fwd_) scoring paths. + + (The disk-resident indexes omit seismic's exact-match fast path, but + still honor the selector per candidate doc.)""" + index = make_index(self.SPEC, corpus) + if residency == "mmap": + index = roundtrip(index, tmp_path / "filtered.idx", nsparse.kUseMmap) + + want_labels, _ = oracle + allowed = np.ascontiguousarray( + np.unique(want_labels[want_labels >= 0])[: K * 5], dtype=np.int32 + ) + assert len(allowed) > K + + selector = nsparse.SetIDSelector(allowed) + p = self.params() + p.set_id_selector(selector) + + _, labels = search(index, queries, params=p) + returned = labels[labels >= 0] + assert np.isin(returned, allowed).all(), "filter must exclude non-members" + + def test_with_id_map(self, corpus, queries, oracle, doc_ids, tmp_path): + """idmap over the index returns the caller's ids, and reloads via mmap + (the delegate's copying read is unsupported, so the whole idmap must be + mmap-loaded).""" + index = make_index(f"idmap,{self.SPEC}", corpus, ids=doc_ids) + index = roundtrip(index, tmp_path / "idmap.idx", nsparse.kUseMmap) + + _, labels = search(index, queries, params=self.params()) + returned = labels[labels >= 0] + assert np.isin(returned, doc_ids).all(), "returned ids must be caller ids" + + want_labels, _ = oracle + want_external = np.where(want_labels >= 0, doc_ids[want_labels], -1) + assert recall_at_k(labels, want_external) >= self.RECALL_FLOOR + + def test_batch_matches_single_query(self, corpus, queries): + """Batched (OpenMP-parallel over queries) results must equal the one-at- + a-time path exactly, or the per-thread scratch is leaking.""" + index = make_index(self.SPEC, corpus) + batch_d, batch_l = search(index, queries, params=self.params()) + single_d, single_l = search_each(index, queries, params=self.params()) + np.testing.assert_array_equal(batch_l, single_l) + np.testing.assert_allclose(batch_d, single_d, rtol=0, atol=0) + + def test_k_larger_than_corpus(self, corpus, queries): + """Short result rows are padded with INVALID_IDX / -1.0, not truncated.""" + small = make_index(self.SPEC, slice_corpus(corpus, 0, 3)) + k = 10 + dists, labels = search(small, queries, k=k, params=self.params()) + assert labels.shape == (queries.n, k) + assert (labels[:, 3:] == PAD_LABEL).all() + assert (dists[:, 3:] == PAD_DIST).all() + + def test_seeded_build_is_reproducible(self, corpus, queries): + """seed= makes the build (and so the search results) reproducible.""" + first = search(make_index(self.SEEDED, corpus), queries, params=self.params()) + second = search(make_index(self.SEEDED, corpus), queries, params=self.params()) + np.testing.assert_array_equal(second[1], first[1]) + np.testing.assert_array_equal(second[0], first[0]) diff --git a/python_tests/test_disk_seismic_index.py b/python_tests/test_disk_seismic_index.py index 85253b3..f5ebbdf 100644 --- a/python_tests/test_disk_seismic_index.py +++ b/python_tests/test_disk_seismic_index.py @@ -7,212 +7,24 @@ """Black-box tests for the disk-resident DiskSeismic index, through SWIG only. -Pins the DiskSeismic-specific contracts: mmap-only reads, the top-k' block -budget, and bit-identical fresh-build vs mmap-reload results. +The disk-resident contract (mmap-only reads, the top-k' block budget, +fresh-build vs mmap-reload parity) is shared with the quantized index and lives +in disk_seismic_contract.py; this suite just binds it to the float spec. """ -import numpy as np -import pytest - import nsparse -from oracle import recall_at_k -from support import ( - K, - PAD_DIST, - PAD_LABEL, - make_index, - roundtrip, - search, - search_each, - slice_corpus, -) +from disk_seismic_contract import DiskSeismicContract SPEC = "disk_seismic,lambda=25|beta=4|alpha=0.4" -# Query knobs. CUT covers every query term (QUERY_NNZ=8); K_PRIME is generous so -# the block budget is not the recall bottleneck. -CUT, K_PRIME = 8, 200 - -# Calibrated against the session corpus; a floor, not a target. -RECALL_FLOOR = 0.80 - -# A seeded spec makes the build deterministic, for tests asserting an exact -# relationship (e.g. strict recall improvement across block budgets). -SEEDED = f"{SPEC}|seed=42" - - -def params(cut=CUT, k_prime=K_PRIME): - return nsparse.DiskSeismicSearchParameters(cut, k_prime) - - -@pytest.fixture(scope="module") -def index(corpus): - return make_index(SPEC, corpus) - - -@pytest.fixture(scope="module") -def mmap_index(corpus, tmp_path_factory): - """Seeded + mmap-reloaded, so block-budget tests exercise the on-disk - (fwd_) path deterministically.""" - index = make_index(SEEDED, corpus) - path = tmp_path_factory.mktemp("disk_seismic") / "seeded.idx" - return roundtrip(index, path, nsparse.kUseMmap) - - -@pytest.mark.parametrize("residency", ["memory", "mmap"]) -def test_happy_case(residency, corpus, queries, oracle, tmp_path): - """factory -> ingest -> build -> query -> accuracy, in both residencies.""" - index = make_index(SPEC, corpus) - assert index.num_vectors() == corpus.n - assert index.get_dimension() == corpus.dim - - if residency == "mmap": - # disk_seismic is mmap-only, so the copying flag (0) would raise; the - # mmap flag is required. The persisted count must survive the reload. - index = roundtrip(index, tmp_path / "disk_seismic.idx", nsparse.kUseMmap) - assert index.num_vectors() == corpus.n - - dists, labels = search(index, queries, params=params()) - assert labels.shape == (queries.n, K) - assert dists.shape == (queries.n, K) - assert (labels[:, 0] >= 0).all(), "every query must return at least one hit" - - want_labels, _ = oracle - assert recall_at_k(labels, want_labels) >= RECALL_FLOOR - - -def test_fresh_build_matches_mmap_reload(index, queries, tmp_path): - """The in-RAM build (vectors_) and the mmap reload (inline forward index) - are two different code paths that must return identical results.""" - p = params() - before_d, before_l = search(index, queries, params=p) - reloaded = roundtrip(index, tmp_path / "disk_seismic.idx", nsparse.kUseMmap) - after_d, after_l = search(reloaded, queries, params=p) - np.testing.assert_array_equal(after_l, before_l) - np.testing.assert_allclose(after_d, before_d, rtol=1e-6, atol=1e-6) - - -def test_copying_read_throws(index, tmp_path): - """disk_seismic is mmap-only: reloading without kUseMmap must raise.""" - path = tmp_path / "disk_seismic.idx" - nsparse.write_index(index, str(path)) - with pytest.raises(RuntimeError, match="mmap-only"): - nsparse.read_index(str(path), 0) - - -@pytest.mark.parametrize("bad", [0, -1]) -def test_rejects_non_positive_block_budget(index, queries, bad): - """A non-positive k' (block budget) is rejected at search time.""" - with pytest.raises(ValueError, match="must be positive"): - search(index, queries, params=params(k_prime=bad)) - - -def test_block_budget_is_monotone(mmap_index, queries, oracle): - """A larger block budget scores a superset of blocks, so recall never - regresses and a big enough budget beats the smallest one. Runs on the - seeded, mmap-reloaded index (deterministic, on-disk block-read path).""" - want_labels, _ = oracle - recalls = [ - recall_at_k( - search(mmap_index, queries, params=params(k_prime=kp))[1], want_labels - ) - for kp in [1, 4, 16, 64, 256] - ] - for lo, hi in zip(recalls, recalls[1:]): - assert hi >= lo - 1e-9, f"recall regressed across k': {recalls}" - assert recalls[-1] > recalls[0], "a larger budget must eventually help" - - -def test_block_budget_saturates(mmap_index, queries): - """Past the candidate-block count, more budget changes nothing.""" - big_d, big_l = search(mmap_index, queries, params=params(k_prime=10**6)) - bigger_d, bigger_l = search(mmap_index, queries, params=params(k_prime=2 * 10**6)) - np.testing.assert_array_equal(bigger_l, big_l) - np.testing.assert_array_equal(bigger_d, big_d) - - -def test_empty_index_roundtrip(queries, tmp_path): - """An un-built, empty index writes, mmap-reloads, and returns all padding.""" - empty = nsparse.index_factory(queries.dim, SPEC) - path = tmp_path / "disk_seismic_empty.idx" - nsparse.write_index(empty, str(path)) - mapped = nsparse.read_index(str(path), nsparse.kUseMmap) - assert mapped.num_vectors() == 0 - dists, labels = search(mapped, queries, params=params()) - assert (labels == PAD_LABEL).all() - assert (dists == PAD_DIST).all() - - -def test_search_before_build(corpus, queries): - """Searching an added-but-unbuilt index yields empty results, not an error. - - Pinned deliberately: it is a silent-empty footgun, so a change here should - be a conscious one.""" - index = nsparse.index_factory(corpus.dim, SPEC) - index.add(corpus.n, corpus.indptr, corpus.indices, corpus.values) - _, labels = search(index, queries, params=params()) - assert (labels == PAD_LABEL).all() - - -def test_filtered_search(index, queries, oracle): - """An id selector larger than k filters results to its members. - - (DiskSeismic omits seismic's exact-match fast path, but still honors the - selector per candidate doc.)""" - want_labels, _ = oracle - allowed = np.ascontiguousarray( - np.unique(want_labels[want_labels >= 0])[: K * 5], dtype=np.int32 - ) - assert len(allowed) > K - - selector = nsparse.SetIDSelector(allowed) - p = params() - p.set_id_selector(selector) - - _, labels = search(index, queries, params=p) - returned = labels[labels >= 0] - assert np.isin(returned, allowed).all(), "filter must exclude non-members" - - -def test_with_id_map(corpus, queries, oracle, doc_ids, tmp_path): - """idmap over disk_seismic returns the caller's ids, and reloads via mmap - (the delegate's copying read is unsupported, so the whole idmap must be - mmap-loaded).""" - index = make_index(f"idmap,{SPEC}", corpus, ids=doc_ids) - index = roundtrip(index, tmp_path / "idmap_disk_seismic.idx", nsparse.kUseMmap) - - _, labels = search(index, queries, params=params()) - returned = labels[labels >= 0] - assert np.isin(returned, doc_ids).all(), "returned ids must be caller ids" - - want_labels, _ = oracle - want_external = np.where(want_labels >= 0, doc_ids[want_labels], -1) - assert recall_at_k(labels, want_external) >= RECALL_FLOOR - - -def test_batch_matches_single_query(index, queries): - """Batched (OpenMP-parallel over queries) results must equal the one-at-a- - time path exactly, or the per-thread dense/visited scratch is leaking.""" - batch_d, batch_l = search(index, queries, params=params()) - single_d, single_l = search_each(index, queries, params=params()) - np.testing.assert_array_equal(batch_l, single_l) - np.testing.assert_allclose(batch_d, single_d, rtol=0, atol=0) - - -def test_k_larger_than_corpus(corpus, queries): - """Short result rows are padded with INVALID_IDX / -1.0, not truncated.""" - small = make_index(SPEC, slice_corpus(corpus, 0, 3)) - k = 10 - dists, labels = search(small, queries, k=k, params=params()) - assert labels.shape == (queries.n, k) - assert (labels[:, 3:] == PAD_LABEL).all() - assert (dists[:, 3:] == PAD_DIST).all() +class TestDiskSeismic(DiskSeismicContract): + SPEC = SPEC + SEEDED = f"{SPEC}|seed=42" + # Calibrated against the session corpus; a floor, not a target. + RECALL_FLOOR = 0.80 -def test_seeded_build_is_reproducible(corpus, queries): - """seed= makes the build (and so the search results) reproducible.""" - seeded = f"{SPEC}|seed=42" - first = search(make_index(seeded, corpus), queries, params=params()) - second = search(make_index(seeded, corpus), queries, params=params()) - np.testing.assert_array_equal(second[1], first[1]) - np.testing.assert_array_equal(second[0], first[0]) + def params(self, cut=None, k_prime=None): + cut = self.CUT if cut is None else cut + k_prime = self.K_PRIME if k_prime is None else k_prime + return nsparse.DiskSeismicSearchParameters(cut, k_prime) diff --git a/python_tests/test_disk_seismic_sq_index.py b/python_tests/test_disk_seismic_sq_index.py new file mode 100644 index 0000000..5a64fe0 --- /dev/null +++ b/python_tests/test_disk_seismic_sq_index.py @@ -0,0 +1,107 @@ +# Copyright OpenSearch Contributors +# SPDX-License-Identifier: Apache-2.0 +# +# The OpenSearch Contributors require contributions made to +# this file be licensed under the Apache-2.0 license or a +# compatible open source license. + +"""Black-box tests for the scalar-quantized DiskSeismic index, through SWIG only. + +The shared disk-resident contract comes from disk_seismic_contract.py (bound +here to the quantized spec); this file adds only what is quantization-specific: +both quantizer widths, the on-disk size win, the query-range override, and the +factory's quantizer handling. +""" + +import os + +import numpy as np +import pytest + +import nsparse +from oracle import recall_at_k +from support import make_index, search +from disk_seismic_contract import DiskSeismicContract + +VMIN, VMAX = 0.0, 1.0 + +# The corpus values span [0.1, 1.0]; the [0, 1] quantizer range covers them. +SPEC = f"disk_seismic_sq,quantizer=8bit|vmin={VMIN}|vmax={VMAX}|lambda=25|beta=4|alpha=0.4" + + +def sq_spec(width): + return f"disk_seismic_sq,quantizer={width}|vmin={VMIN}|vmax={VMAX}|lambda=25|beta=4|alpha=0.4" + + +class TestDiskSeismicSQ(DiskSeismicContract): + SPEC = SPEC + SEEDED = f"{SPEC}|seed=42" + # A floor with headroom, not a target: the build RNG is unseeded and 8-bit + # quantization adds noise, so a tight bound would flake. + RECALL_FLOOR = 0.75 + + def params(self, cut=None, k_prime=None): + cut = self.CUT if cut is None else cut + k_prime = self.K_PRIME if k_prime is None else k_prime + return nsparse.DiskSeismicSQSearchParameters(VMIN, VMAX, cut, k_prime) + + # --- quantization-specific --- + + @pytest.mark.parametrize("width", ["8bit", "16bit"]) + def test_quantizer_widths(self, corpus, queries, oracle, width): + """Both quantizer widths clear the recall floor.""" + idx = make_index(sq_spec(width), corpus) + _, labels = search(idx, queries, params=self.params()) + want_labels, _ = oracle + assert recall_at_k(labels, want_labels) >= self.RECALL_FLOOR + + @pytest.mark.parametrize("width", ["8bit", "16bit"]) + def test_smaller_than_plain_disk_seismic(self, corpus, tmp_path, width): + """The whole point of quantization: codes shrink the on-disk index + versus the float disk_seismic, for the same corpus and cluster params. + Both widths must win.""" + float_index = make_index("disk_seismic,lambda=25|beta=4|alpha=0.4|seed=42", corpus) + float_path = tmp_path / "disk_seismic_float.idx" + nsparse.write_index(float_index, str(float_path)) + + sq_index = make_index(f"{sq_spec(width)}|seed=42", corpus) + sq_path = tmp_path / f"disk_seismic_sq_{width}.idx" + nsparse.write_index(sq_index, str(sq_path)) + + assert os.path.getsize(sq_path) < os.path.getsize(float_path) + + def test_query_range_override(self, corpus, queries): + """A DiskSeismicSQSearchParameters with a non-index range re-encodes the + query at that range (the query_quantizer override). The corpus values + span [0.1, 1.0]; vmax=0.5 re-scales/clips the query, so it must actually + change the decoded scores versus the default range -- if query_quantizer + ignored the override this would silently no-op.""" + index = make_index(self.SPEC, corpus) + default_d, _ = search(index, queries, params=self.params()) + + clipped = nsparse.DiskSeismicSQSearchParameters( + 0.0, 0.5, self.CUT, self.K_PRIME + ) + clipped_d, clipped_l = search(index, queries, params=clipped) + + assert clipped_l.shape == default_d.shape + assert clipped_l[clipped_l >= 0].size > 0, "clipping must not empty every result" + assert not np.allclose(clipped_d, default_d), "query-range override had no effect" + + def test_factory_downcast_exposes_quantizer(self, corpus): + """DSSQ is in the %factory downcast lists, so index_factory returns the + concrete type and the DSSQ-only get_scalar_quantizer() is reachable (it + is not a base Index method). The two widths must report distinct types.""" + + def qtype(width): + idx = nsparse.index_factory(corpus.dim, sq_spec(width)) + return idx.get_scalar_quantizer().get_quantizer_type() + + assert qtype("8bit") != qtype("16bit"), "widths must report distinct types" + + @pytest.mark.parametrize("bad", ["4bit", "int8", ""]) + def test_factory_rejects_unknown_quantizer(self, corpus, bad): + """An unrecognized quantizer= is rejected, not silently treated as 8bit.""" + spec = f"disk_seismic_sq,quantizer={bad}|vmin={VMIN}|vmax={VMAX}|lambda=25" + with pytest.raises(ValueError, match="quantizer"): + nsparse.index_factory(corpus.dim, spec) diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index d404ee3..22ecc89 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -31,6 +31,7 @@ set(NSPARSE_TEST_SRC scalar_quantizer_test.cpp seismic_common_test.cpp disk_seismic_index_test.cpp + disk_seismic_scalar_quantized_index_test.cpp seismic_index_test.cpp seismic_invlists_writer_test.cpp seismic_scalar_quantized_index_test.cpp diff --git a/tests/disk_seismic_index_test.cpp b/tests/disk_seismic_index_test.cpp index 7b6f9e9..373086f 100644 --- a/tests/disk_seismic_index_test.cpp +++ b/tests/disk_seismic_index_test.cpp @@ -11,13 +11,8 @@ #include -#include #include -#include #include -#include -#include -#include #include #include "nsparse/index_factory.h" @@ -25,112 +20,10 @@ #include "nsparse/io/index_io.h" #include "nsparse/seismic_index.h" #include "nsparse/types.h" +#include "tests/disk_seismic_test_util.h" namespace nsparse { -namespace { - -// Fixed cluster params + seed so SeismicIndex and DiskSeismicIndex build -// bit-identical clusters and their results can be compared exactly. -constexpr int kLambda = 32; -constexpr int kBeta = 8; -constexpr float kAlpha = 0.4F; -constexpr int kSeed = 42; -constexpr int kDim = 400; - -SeismicClusterParameters cluster_params() { - return {.lambda = kLambda, .beta = kBeta, .alpha = kAlpha, .seed = kSeed}; -} - -struct CSR { - std::vector indptr; - std::vector indices; - std::vector values; - idx_t n = 0; -}; - -// A reproducible random sparse corpus: each row has a random number of distinct -// terms (sorted ascending, CSR convention) with values in (0, 1]. -CSR make_corpus(idx_t rows, unsigned seed) { - std::mt19937 rng(seed); - std::uniform_int_distribution nnz_dist(5, 25); - std::uniform_int_distribution term_dist(0, kDim - 1); - std::uniform_real_distribution val_dist(0.05F, 1.0F); - CSR c; - c.n = rows; - c.indptr.push_back(0); - for (idx_t r = 0; r < rows; ++r) { - std::set terms; - const int nnz = nnz_dist(rng); - while (static_cast(terms.size()) < nnz) - terms.insert(term_dist(rng)); - for (const int t : terms) { - c.indices.push_back(static_cast(t)); - c.values.push_back(val_dist(rng)); - } - c.indptr.push_back(static_cast(c.indices.size())); - } - return c; -} - -void add_corpus(Index& index, const CSR& c) { - index.add(c.n, c.indptr.data(), c.indices.data(), c.values.data()); -} - -// Index file removed on destruction. write_index/read_index take char*. -class TempIndexFile { -public: - explicit TempIndexFile(const std::string& name) - : path_(std::filesystem::temp_directory_path() / name) { - std::filesystem::remove(path_); - } - ~TempIndexFile() { std::filesystem::remove(path_); } - TempIndexFile(const TempIndexFile&) = delete; - TempIndexFile& operator=(const TempIndexFile&) = delete; - char* c_str() { return path_str_.data(); } - -private: - std::filesystem::path path_{}; - std::string path_str_ = path_.string(); -}; - -using ScoreIds = - std::pair>, std::vector>>; - -// K' large enough to select every candidate block (saturates the budget). -constexpr int kAllBlocks = 1'000'000; - -ScoreIds search_all(Index& index, const CSR& queries, int k, - SearchParameters* params) { - std::vector distances(static_cast(queries.n) * k); - std::vector labels(static_cast(queries.n) * k); - index.search(queries.n, queries.indptr.data(), queries.indices.data(), - queries.values.data(), k, distances.data(), labels.data(), - params); - ScoreIds out; - out.first.resize(queries.n); - out.second.resize(queries.n); - for (idx_t q = 0; q < queries.n; ++q) { - for (int j = 0; j < k; ++j) { - out.first[q].push_back(distances[static_cast(q) * k + j]); - out.second[q].push_back(labels[static_cast(q) * k + j]); - } - } - return out; -} - -void expect_same_results(const ScoreIds& a, const ScoreIds& b) { - ASSERT_EQ(a.second.size(), b.second.size()); - for (size_t q = 0; q < a.second.size(); ++q) { - EXPECT_EQ(a.second[q], b.second[q]) << "labels differ at query " << q; - ASSERT_EQ(a.first[q].size(), b.first[q].size()); - for (size_t j = 0; j < a.first[q].size(); ++j) { - EXPECT_FLOAT_EQ(a.first[q][j], b.first[q][j]) - << "score differs at query " << q << " rank " << j; - } - } -} - -} // namespace +using namespace disk_seismic_test; // NOLINT(build/namespaces) // Bit-exact anchor: reading a selected block's vectors from the inline mapping // (fwd_) gives the same result as reading them from the in-RAM CSR (vectors_) diff --git a/tests/disk_seismic_scalar_quantized_index_test.cpp b/tests/disk_seismic_scalar_quantized_index_test.cpp new file mode 100644 index 0000000..e80222c --- /dev/null +++ b/tests/disk_seismic_scalar_quantized_index_test.cpp @@ -0,0 +1,334 @@ +/** + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +#include "nsparse/disk_seismic_scalar_quantized_index.h" + +#include + +#include +#include +#include +#include +#include +#include +#include +#include + +#include "nsparse/disk_seismic_index.h" +#include "nsparse/index_factory.h" +#include "nsparse/io/buffered_io.h" +#include "nsparse/io/index_io.h" +#include "nsparse/types.h" +#include "nsparse/utils/scalar_quantizer.h" +#include "tests/disk_seismic_test_util.h" + +namespace nsparse { +using namespace disk_seismic_test; // NOLINT(build/namespaces) +namespace { + +// The quantizer header (type u8 + vmin f32 + vmax f32) is the first thing +// write_index writes after the common IndexHeader, so the type byte sits right +// at kIndexHeaderSize. +constexpr size_t kQuantizerTypeOffset = kIndexHeaderSize; + +template +T read_field(const std::string& path, size_t offset) { + std::ifstream in(path, std::ios::binary); + in.seekg(static_cast(offset)); + T value{}; + in.read(reinterpret_cast(&value), sizeof(T)); + return value; +} + +// A load-time guard rejects what the writer cannot produce, so the only way to +// reach one is to overwrite a field of an otherwise valid file in place. +template +void patch_field(const std::string& path, size_t offset, T value) { + std::fstream out(path, std::ios::binary | std::ios::in | std::ios::out); + out.seekp(static_cast(offset)); + out.write(reinterpret_cast(&value), sizeof(T)); +} + +// disk_seismic_sq is mmap-only, so the copying read throws "mmap-only" before +// any quantizer guard runs; only the mmap read reaches the guard. The message +// is checked, not just the type, so a guard silently dropped (and something +// downstream throwing instead) does not read as a pass. +void expect_mmap_read_rejected(char* path, const char* fragment) { + try { + std::unique_ptr loaded(read_index(path, IndexIoFlag::kUseMmap)); + ADD_FAILURE() << "accepted the corrupted file"; + } catch (const std::runtime_error& error) { + EXPECT_NE(std::string(error.what()).find(fragment), std::string::npos) + << error.what(); + } +} + +// Mean fraction of `reference`'s valid labels recovered by `got`, per query. +double recall(const ScoreIds& got, const ScoreIds& reference) { + double sum = 0.0; + const size_t nq = reference.second.size(); + for (size_t q = 0; q < nq; ++q) { + std::unordered_set hits(got.second[q].begin(), + got.second[q].end()); + size_t found = 0; + size_t total = 0; + for (const idx_t label : reference.second[q]) { + if (label == detail::INVALID_IDX) { + continue; + } + ++total; + if (hits.count(label) > 0) { + ++found; + } + } + if (total > 0) { + sum += static_cast(found) / static_cast(total); + } + } + return nq > 0 ? sum / static_cast(nq) : 1.0; +} + +} // namespace + +// Bit-exact anchor: reading a selected block's codes from the inline mapping +// (fwd_) gives the same result as reading them from the in-RAM CSR (vectors_) +// of a fresh build, at the same K' -- same clusters, same blocks, same codes +// and integer dots, same decode. +TEST(DiskSeismicSQIndex, MappedReloadMatchesFreshBuild8bit) { + const CSR corpus = make_corpus(1500, /*seed=*/1); + const CSR queries = make_corpus(40, /*seed=*/2); + + DiskSeismicScalarQuantizedIndex disk(QuantizerType::QT_8bit, 0.0F, 1.0F, + cluster_params(), kDim); + add_corpus(disk, corpus); + disk.build(); + DiskSeismicSearchParameters params(/*cut=*/25, /*k_prime=*/32); + const ScoreIds fresh = search_all(disk, queries, 10, ¶ms); // vectors_ + + TempIndexFile file("nsparse_disk_seismic_sq_parity8.idx"); + write_index(&disk, file.c_str()); + std::unique_ptr mapped( + read_index(file.c_str(), IndexIoFlag::kUseMmap)); + ASSERT_NE(mapped, nullptr); + EXPECT_EQ(mapped->num_vectors(), static_cast(corpus.n)); + expect_same_results(search_all(*mapped, queries, 10, ¶ms), + fresh); // fwd_ +} + +// Same fresh-vs-mapped parity at 16-bit (the other quantizer width). +TEST(DiskSeismicSQIndex, MappedReloadMatchesFreshBuild16bit) { + const CSR corpus = make_corpus(1500, /*seed=*/1); + const CSR queries = make_corpus(40, /*seed=*/2); + + DiskSeismicScalarQuantizedIndex disk(QuantizerType::QT_16bit, 0.0F, 1.0F, + cluster_params(), kDim); + add_corpus(disk, corpus); + disk.build(); + DiskSeismicSearchParameters params(/*cut=*/25, /*k_prime=*/32); + const ScoreIds fresh = search_all(disk, queries, 10, ¶ms); + + TempIndexFile file("nsparse_disk_seismic_sq_parity16.idx"); + write_index(&disk, file.c_str()); + std::unique_ptr mapped( + read_index(file.c_str(), IndexIoFlag::kUseMmap)); + ASSERT_NE(mapped, nullptr); + expect_same_results(search_all(*mapped, queries, 10, ¶ms), fresh); +} + +// add() encodes to the quantizer's element width (1 byte for 8-bit, 2 for +// 16-bit), so the forward index and summaries are stored as codes. +TEST(DiskSeismicSQIndex, AddEncodesToCodeWidth) { + const CSR corpus = make_corpus(50, /*seed=*/9); + DiskSeismicScalarQuantizedIndex eight(QuantizerType::QT_8bit, 0.0F, 1.0F, + cluster_params(), kDim); + add_corpus(eight, corpus); + ASSERT_NE(eight.get_vectors(), nullptr); + EXPECT_EQ(eight.get_vectors()->get_element_size(), U8); + + DiskSeismicScalarQuantizedIndex sixteen(QuantizerType::QT_16bit, 0.0F, 1.0F, + cluster_params(), kDim); + add_corpus(sixteen, corpus); + ASSERT_NE(sixteen.get_vectors(), nullptr); + EXPECT_EQ(sixteen.get_vectors()->get_element_size(), U16); +} + +// With the cut covering every query term and the budget saturated, both indexes +// score the identical candidate set, so any ranking difference is quantization +// alone. 8-bit over the corpus's [0.05, 1] range recovers nearly all of float's +// top-k. The floor is conservative, not a target. +TEST(DiskSeismicSQIndex, RecallCloseToFloat) { + const CSR corpus = make_corpus(1500, /*seed=*/11); + const CSR queries = make_corpus(40, /*seed=*/12); + DiskSeismicSearchParameters params(/*cut=*/25, kAllBlocks); + + DiskSeismicIndex fdisk(kDim, cluster_params()); + add_corpus(fdisk, corpus); + fdisk.build(); + const ScoreIds fref = search_all(fdisk, queries, 10, ¶ms); + + DiskSeismicScalarQuantizedIndex qdisk(QuantizerType::QT_8bit, 0.0F, 1.0F, + cluster_params(), kDim); + add_corpus(qdisk, corpus); + qdisk.build(); + const ScoreIds qres = search_all(qdisk, queries, 10, ¶ms); + + EXPECT_GE(recall(qres, fref), 0.80); +} + +// The whole point: codes shrink the on-disk index versus float (forward vals +// 4B/nnz -> 1B or 2B/nnz, summaries likewise), for the same corpus and +// clusters. Both widths must win, so a 16-bit-specific layout or padding +// regression cannot inflate the file past float unnoticed. +TEST(DiskSeismicSQIndex, SmallerThanFloat) { + const CSR corpus = make_corpus(1500, /*seed=*/13); + + DiskSeismicIndex fdisk(kDim, cluster_params()); + add_corpus(fdisk, corpus); + fdisk.build(); + TempIndexFile ffile("nsparse_disk_seismic_sq_float.idx"); + write_index(&fdisk, ffile.c_str()); + const auto float_size = ffile.size(); + + for (const QuantizerType qt : + {QuantizerType::QT_8bit, QuantizerType::QT_16bit}) { + DiskSeismicScalarQuantizedIndex qdisk(qt, 0.0F, 1.0F, cluster_params(), + kDim); + add_corpus(qdisk, corpus); + qdisk.build(); + TempIndexFile qfile("nsparse_disk_seismic_sq_codes.idx"); + write_index(&qdisk, qfile.c_str()); + EXPECT_LT(qfile.size(), float_size) + << "quantized index (" << qfile.size() + << ") not smaller than float (" << float_size << ") at width " + << (qt == QuantizerType::QT_8bit ? 1 : 2); + } +} + +// K' (block budget) must be positive. +TEST(DiskSeismicSQIndex, RejectsNonPositiveBlockBudget) { + const CSR corpus = make_corpus(300, /*seed=*/5); + const CSR queries = make_corpus(2, /*seed=*/6); + DiskSeismicScalarQuantizedIndex disk(QuantizerType::QT_8bit, 0.0F, 1.0F, + cluster_params(), kDim); + add_corpus(disk, corpus); + disk.build(); + std::vector distances(2 * 10); + std::vector labels(2 * 10); + Index& idx = disk; + for (const int bad : {0, -1}) { + DiskSeismicSearchParameters params(/*cut=*/10, bad); + EXPECT_THROW(idx.search(queries.n, queries.indptr.data(), + queries.indices.data(), queries.values.data(), + 10, distances.data(), labels.data(), ¶ms), + std::invalid_argument); + } +} + +// mmap-only: the copying read path is unsupported. +TEST(DiskSeismicSQIndex, CopyingReadThrows) { + const CSR corpus = make_corpus(300, /*seed=*/5); + DiskSeismicScalarQuantizedIndex disk(QuantizerType::QT_8bit, 0.0F, 1.0F, + cluster_params(), kDim); + add_corpus(disk, corpus); + disk.build(); + TempIndexFile file("nsparse_disk_seismic_sq_copying.idx"); + write_index(&disk, file.c_str()); + EXPECT_THROW(read_index(file.c_str(), /*io_flags=*/0), std::runtime_error); + // A stream reader cannot be mapped either, so the flag still throws. + BufferedIOWriter writer; + write_index(&disk, &writer); + BufferedIOReader reader(writer.data()); + EXPECT_THROW(read_index(&reader, IndexIoFlag::kUseMmap), + std::runtime_error); +} + +// bytes_per_value() reads anything that is not QT_8bit as 16-bit, so an +// undefined type would pick an element width instead of being rejected, and the +// codes behind it would be strided at that width. The mmap reader parses the +// quantizer header itself, so the guard has to hold there. +TEST(DiskSeismicSQIndex, MmapReadRejectsUnknownQuantizerType) { + const CSR corpus = make_corpus(300, /*seed=*/21); + DiskSeismicScalarQuantizedIndex disk(QuantizerType::QT_8bit, 0.0F, 1.0F, + cluster_params(), kDim); + add_corpus(disk, corpus); + disk.build(); + TempIndexFile file("nsparse_disk_seismic_sq_unknown_type.idx"); + write_index(&disk, file.c_str()); + + // Fails loudly if the layout moves, rather than patching some other field. + ASSERT_EQ(read_field(file.c_str(), kQuantizerTypeOffset), + QuantizerType::QT_8bit); + patch_field(std::string(file.c_str()), kQuantizerTypeOffset, + static_cast(2)); + + expect_mmap_read_rejected(file.c_str(), "unknown quantizer type"); +} + +// A type that is a valid enum but not the width the codes were stored at: +// search would stride the stored codes at the quantizer's width and read past +// the end of the array, or halfway into each code. +TEST(DiskSeismicSQIndex, MmapReadRejectsWidthMismatch) { + const CSR corpus = make_corpus(300, /*seed=*/22); + DiskSeismicScalarQuantizedIndex disk(QuantizerType::QT_8bit, 0.0F, 1.0F, + cluster_params(), kDim); + add_corpus(disk, corpus); + disk.build(); + TempIndexFile file("nsparse_disk_seismic_sq_width_mismatch.idx"); + write_index(&disk, file.c_str()); + + ASSERT_EQ(read_field(file.c_str(), kQuantizerTypeOffset), + QuantizerType::QT_8bit); + // Valid enum, but the codes were written at 8-bit width. + patch_field(std::string(file.c_str()), kQuantizerTypeOffset, + QuantizerType::QT_16bit); + + expect_mmap_read_rejected(file.c_str(), + "element size disagrees with its quantizer type"); +} + +// An empty index (no docs) round-trips: count 0, no results, no crash. The +// quantizer header still writes and reads back. +TEST(DiskSeismicSQIndex, EmptyIndexRoundTrip) { + DiskSeismicScalarQuantizedIndex disk(QuantizerType::QT_8bit, 0.0F, 1.0F, + cluster_params(), kDim); + TempIndexFile file("nsparse_disk_seismic_sq_empty.idx"); + write_index(&disk, file.c_str()); + std::unique_ptr mapped( + read_index(file.c_str(), IndexIoFlag::kUseMmap)); + ASSERT_NE(mapped, nullptr); + EXPECT_EQ(mapped->num_vectors(), 0U); + + const CSR queries = make_corpus(3, /*seed=*/6); + std::vector distances(3 * 5, -1.0F); + std::vector labels(3 * 5, detail::INVALID_IDX); + DiskSeismicSearchParameters params(10, 50); + mapped->search(queries.n, queries.indptr.data(), queries.indices.data(), + queries.values.data(), 5, distances.data(), labels.data(), + ¶ms); // must not crash; nothing found +} + +// index_factory understands the "disk_seismic_sq" descriptor and its quantizer +// parameter. +TEST(DiskSeismicSQIndex, FactoryCreatesIt) { + std::unique_ptr eight(index_factory( + kDim, "disk_seismic_sq,quantizer=8bit|lambda=32|beta=8|seed=42")); + ASSERT_NE(eight, nullptr); + EXPECT_EQ(eight->id(), DiskSeismicScalarQuantizedIndex::name); + + std::unique_ptr sixteen(index_factory( + kDim, "disk_seismic_sq,quantizer=16bit|lambda=32|beta=8|seed=42")); + ASSERT_NE(sixteen, nullptr); + auto* typed = + dynamic_cast(sixteen.get()); + ASSERT_NE(typed, nullptr); + EXPECT_EQ(typed->get_scalar_quantizer().get_quantizer_type(), + QuantizerType::QT_16bit); +} + +} // namespace nsparse diff --git a/tests/disk_seismic_test_util.h b/tests/disk_seismic_test_util.h new file mode 100644 index 0000000..5528709 --- /dev/null +++ b/tests/disk_seismic_test_util.h @@ -0,0 +1,139 @@ +/** + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +#ifndef DISK_SEISMIC_TEST_UTIL_H +#define DISK_SEISMIC_TEST_UTIL_H + +#include + +#include +#include +#include +#include +#include +#include +#include + +#include "nsparse/index.h" +#include "nsparse/io/index_io.h" +#include "nsparse/seismic_index.h" +#include "nsparse/types.h" + +// Shared fixtures for the two disk-resident index test suites (DiskSeismicIndex +// and DiskSeismicScalarQuantizedIndex): the same random corpus, temp file, and +// batch-search harness, so each suite adds only its own assertions. +namespace nsparse::disk_seismic_test { + +// Fixed cluster params + seed so builds are reproducible and (where both +// indexes are compared) they select the same candidate blocks. +inline constexpr int kLambda = 32; +inline constexpr int kBeta = 8; +inline constexpr float kAlpha = 0.4F; +inline constexpr int kSeed = 42; +inline constexpr int kDim = 400; + +// K' large enough to select every candidate block (saturates the budget). +inline constexpr int kAllBlocks = 1'000'000; + +inline SeismicClusterParameters cluster_params() { + return {.lambda = kLambda, .beta = kBeta, .alpha = kAlpha, .seed = kSeed}; +} + +struct CSR { + std::vector indptr; + std::vector indices; + std::vector values; + idx_t n = 0; +}; + +// A reproducible random sparse corpus: each row has a random number of distinct +// terms (sorted ascending, CSR convention) with values in (0, 1] -- so the +// default quantizer range [0, 1] covers them without clipping. +inline CSR make_corpus(idx_t rows, unsigned seed) { + std::mt19937 rng(seed); + std::uniform_int_distribution nnz_dist(5, 25); + std::uniform_int_distribution term_dist(0, kDim - 1); + std::uniform_real_distribution val_dist(0.05F, 1.0F); + CSR c; + c.n = rows; + c.indptr.push_back(0); + for (idx_t r = 0; r < rows; ++r) { + std::set terms; + const int nnz = nnz_dist(rng); + while (static_cast(terms.size()) < nnz) { + terms.insert(term_dist(rng)); + } + for (const int t : terms) { + c.indices.push_back(static_cast(t)); + c.values.push_back(val_dist(rng)); + } + c.indptr.push_back(static_cast(c.indices.size())); + } + return c; +} + +inline void add_corpus(Index& index, const CSR& c) { + index.add(c.n, c.indptr.data(), c.indices.data(), c.values.data()); +} + +// Index file removed on destruction. write_index/read_index take char*. +class TempIndexFile { +public: + explicit TempIndexFile(const std::string& name) + : path_(std::filesystem::temp_directory_path() / name) { + std::filesystem::remove(path_); + } + ~TempIndexFile() { std::filesystem::remove(path_); } + TempIndexFile(const TempIndexFile&) = delete; + TempIndexFile& operator=(const TempIndexFile&) = delete; + char* c_str() { return path_str_.data(); } + std::uintmax_t size() const { return std::filesystem::file_size(path_); } + +private: + std::filesystem::path path_{}; + std::string path_str_ = path_.string(); +}; + +using ScoreIds = + std::pair>, std::vector>>; + +inline ScoreIds search_all(Index& index, const CSR& queries, int k, + SearchParameters* params) { + std::vector distances(static_cast(queries.n) * k); + std::vector labels(static_cast(queries.n) * k); + index.search(queries.n, queries.indptr.data(), queries.indices.data(), + queries.values.data(), k, distances.data(), labels.data(), + params); + ScoreIds out; + out.first.resize(queries.n); + out.second.resize(queries.n); + for (idx_t q = 0; q < queries.n; ++q) { + for (int j = 0; j < k; ++j) { + out.first[q].push_back(distances[static_cast(q) * k + j]); + out.second[q].push_back(labels[static_cast(q) * k + j]); + } + } + return out; +} + +inline void expect_same_results(const ScoreIds& a, const ScoreIds& b) { + ASSERT_EQ(a.second.size(), b.second.size()); + for (size_t q = 0; q < a.second.size(); ++q) { + EXPECT_EQ(a.second[q], b.second[q]) << "labels differ at query " << q; + ASSERT_EQ(a.first[q].size(), b.first[q].size()); + for (size_t j = 0; j < a.first[q].size(); ++j) { + EXPECT_FLOAT_EQ(a.first[q][j], b.first[q][j]) + << "score differs at query " << q << " rank " << j; + } + } +} + +} // namespace nsparse::disk_seismic_test + +#endif // DISK_SEISMIC_TEST_UTIL_H diff --git a/tests/seismic_scalar_quantized_index_test.cpp b/tests/seismic_scalar_quantized_index_test.cpp index 9570f7e..8202df8 100644 --- a/tests/seismic_scalar_quantized_index_test.cpp +++ b/tests/seismic_scalar_quantized_index_test.cpp @@ -998,7 +998,7 @@ class ScopedMmapAdvise { #endif // Where read_index leaves off before the payload. SESQ's payload opens with the -// quantizer header write_quantizer_header wrote. +// quantization header write_quantization_header wrote. constexpr size_t kQuantizerTypeOffset = kIndexHeaderSize; constexpr size_t kVminOffset = kQuantizerTypeOffset + sizeof(QuantizerType);