diff --git a/nsparse/disk_seismic_index_base.h b/nsparse/disk_seismic_index_base.h index 8af4dc0..be2f53f 100644 --- a/nsparse/disk_seismic_index_base.h +++ b/nsparse/disk_seismic_index_base.h @@ -44,6 +44,21 @@ class DiskSeismicIndexBase : public MmapIndex, public IndexIO { const float* values) override; void build() override; + // read_mcsr (the kMmap path) borrows vectors_ from the mapping but never + // touches num_vectors_, which otherwise only add() maintains; write_index + // and search read the member directly, so an out-of-sync count would + // persist nv=0 and gate search to empty. Sync it here after the base read. + // (kInMemory routes through add(), which already set it, so re-reading + // get_vectors() is then a harmless no-op.) + void read_csr(const char* file_path, + Residency residency = Residency::kInMemory) override { + MmapIndex::read_csr(file_path, residency); + const auto* v = get_vectors(); + if (v != nullptr) { + num_vectors_ = v->num_vectors(); + } + } + protected: DiskSeismicIndexBase(int dim, SeismicClusterParameters parameter); diff --git a/tests/csr_interchange_test_util.h b/tests/csr_interchange_test_util.h new file mode 100644 index 0000000..034d64f --- /dev/null +++ b/tests/csr_interchange_test_util.h @@ -0,0 +1,80 @@ +/** + * 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 CSR_INTERCHANGE_TEST_UTIL_H +#define CSR_INTERCHANGE_TEST_UTIL_H + +#include +#include +#include +#include +#include +#include +#include + +// Shared helpers for the mmap-CSR build path, used by both the regular and the +// disk-resident index suites: write a corpus as an interchange CSR (the layout +// csr_layout::convert consumes) and manage the interchange + native temp files. +namespace nsparse::csr_test { + +// Writes a corpus as an interchange CSR: int64 header {rows, num_cols, nnz}, +// int64 indptr[rows + 1], int32 indices[nnz], float values[nnz]. Templated on +// the corpus struct (any type exposing .n / .indptr / .indices / .values), so it +// serves any test corpus. The values are written verbatim, so a convert + +// read_csr(kMmap) build sees the exact same vectors as add(). +template +void write_interchange_csr(const std::string& path, const Corpus& c, + int num_cols) { + std::ofstream out(path, std::ios::binary); + const std::array header = { + static_cast(c.n), static_cast(num_cols), + static_cast(c.indices.size())}; + out.write(reinterpret_cast(header.data()), + header.size() * sizeof(int64_t)); + const std::vector indptr64(c.indptr.begin(), c.indptr.end()); + out.write(reinterpret_cast(indptr64.data()), + static_cast(indptr64.size() * sizeof(int64_t))); + const std::vector indices32(c.indices.begin(), c.indices.end()); + out.write(reinterpret_cast(indices32.data()), + static_cast(indices32.size() * sizeof(int32_t))); + out.write(reinterpret_cast(c.values.data()), + static_cast(c.values.size() * sizeof(float))); +} + +// An interchange CSR temp file and the native path convert writes it to, both +// removed on destruction. +class TempCsrFiles { +public: + explicit TempCsrFiles(const std::string& stem) + : interchange_(std::filesystem::temp_directory_path() / (stem + ".csr")), + native_(std::filesystem::temp_directory_path() / (stem + ".mcsr")) { + std::error_code ignored; + std::filesystem::remove(interchange_, ignored); + std::filesystem::remove(native_, ignored); + } + ~TempCsrFiles() { + std::error_code ignored; + std::filesystem::remove(interchange_, ignored); + std::filesystem::remove(native_, ignored); + } + TempCsrFiles(const TempCsrFiles&) = delete; + TempCsrFiles& operator=(const TempCsrFiles&) = delete; + const std::string& interchange() const { return interchange_str_; } + const std::string& native() const { return native_str_; } + +private: + std::filesystem::path interchange_; + std::filesystem::path native_; + std::string interchange_str_ = interchange_.string(); + std::string native_str_ = native_.string(); +}; + +} // namespace nsparse::csr_test + +#endif // CSR_INTERCHANGE_TEST_UTIL_H diff --git a/tests/disk_seismic_index_test.cpp b/tests/disk_seismic_index_test.cpp index 373086f..d2833b9 100644 --- a/tests/disk_seismic_index_test.cpp +++ b/tests/disk_seismic_index_test.cpp @@ -20,6 +20,7 @@ #include "nsparse/io/index_io.h" #include "nsparse/seismic_index.h" #include "nsparse/types.h" +#include "nsparse/utils/csr_layout.h" #include "tests/disk_seismic_test_util.h" namespace nsparse { @@ -49,6 +50,48 @@ TEST(DiskSeismicIndex, MappedReloadMatchesFreshBuild) { fresh); // fwd_ } +// Building from a native CSR borrowed via mmap must match building the same +// corpus fed through add(): convert -> read_csr(kMmap) -> build -> persist -> +// mmap-reload -> search is bit-exact to the add()-fed build. This is also the +// regression guard for the num_vectors_ desync -- read_csr(kMmap) borrows +// vectors_ without running add(), so unless num_vectors_ is synced the written +// index records nv=0 and a reloaded search returns nothing. +TEST(DiskSeismicIndex, MmapCsrBuildMatchesAddBuild) { + const CSR corpus = make_corpus(1500, /*seed=*/1); + const CSR queries = make_corpus(40, /*seed=*/2); + DiskSeismicSearchParameters params(/*cut=*/10, /*k_prime=*/32); + + // Reference: the same corpus fed through add(), then reloaded via mmap. + DiskSeismicIndex added(kDim, cluster_params()); + add_corpus(added, corpus); + added.build(); + TempIndexFile added_file("nsparse_disk_seismic_addbuild.idx"); + write_index(&added, added_file.c_str()); + std::unique_ptr added_mapped( + read_index(added_file.c_str(), IndexIoFlag::kUseMmap)); + ASSERT_NE(added_mapped, nullptr); + const ScoreIds fresh = search_all(*added_mapped, queries, 10, ¶ms); + + // Under test: build from a native CSR of the same corpus, borrowed via mmap. + TempCsrFiles csr("nsparse_disk_seismic_source"); + write_interchange_csr(csr.interchange(), corpus, kDim); + csr_layout::convert(csr.interchange(), csr.native()); + + DiskSeismicIndex mapped_build(kDim, cluster_params()); + mapped_build.read_csr(csr.native().c_str(), Residency::kMmap); + // num_vectors_ must reflect the borrowed vectors (regression guard). + ASSERT_EQ(mapped_build.num_vectors(), static_cast(corpus.n)); + mapped_build.build(); + TempIndexFile built_file("nsparse_disk_seismic_mmapbuild.idx"); + write_index(&mapped_build, built_file.c_str()); + std::unique_ptr built_mapped( + read_index(built_file.c_str(), IndexIoFlag::kUseMmap)); + ASSERT_NE(built_mapped, nullptr); + // The persisted doc count must survive the mmap-CSR build. + EXPECT_EQ(built_mapped->num_vectors(), static_cast(corpus.n)); + expect_same_results(search_all(*built_mapped, queries, 10, ¶ms), fresh); +} + // K' is a monotonic depth budget: the global top-K' blocks are nested by // summary score, so a larger budget scores a superset of docs and every rank's // score can only improve. This also proves the parameter is actually consumed. diff --git a/tests/disk_seismic_test_util.h b/tests/disk_seismic_test_util.h index 5528709..9835c9d 100644 --- a/tests/disk_seismic_test_util.h +++ b/tests/disk_seismic_test_util.h @@ -24,6 +24,7 @@ #include "nsparse/io/index_io.h" #include "nsparse/seismic_index.h" #include "nsparse/types.h" +#include "tests/csr_interchange_test_util.h" // Shared fixtures for the two disk-resident index test suites (DiskSeismicIndex // and DiskSeismicScalarQuantizedIndex): the same random corpus, temp file, and @@ -82,6 +83,11 @@ inline void add_corpus(Index& index, const CSR& c) { index.add(c.n, c.indptr.data(), c.indices.data(), c.values.data()); } +// The mmap-CSR build helpers are generic (see csr_interchange_test_util.h); +// re-exported here so the disk suite reaches them via `disk_seismic_test`. +using csr_test::TempCsrFiles; +using csr_test::write_interchange_csr; + // Index file removed on destruction. write_index/read_index take char*. class TempIndexFile { public: diff --git a/tests/seismic_index_test.cpp b/tests/seismic_index_test.cpp index d09168c..35a9fd9 100644 --- a/tests/seismic_index_test.cpp +++ b/tests/seismic_index_test.cpp @@ -17,6 +17,7 @@ #include #include #include +#include #include #include #include @@ -28,6 +29,8 @@ #include "nsparse/io/index_io.h" #include "nsparse/sparse_vectors.h" #include "nsparse/types.h" +#include "nsparse/utils/csr_layout.h" +#include "tests/csr_interchange_test_util.h" namespace nsparse { namespace { @@ -956,6 +959,50 @@ void expect_both_reads_rejected(char* path, const char* fragment) { } } +// Raw-array CSR corpus. A reproducible random corpus with distinct, ascending +// terms per row (CSR convention) and values in (0, 1]. +struct RawCsr { + std::vector indptr; + std::vector indices; + std::vector values; + idx_t n = 0; +}; + +RawCsr make_raw_corpus(idx_t rows, int dim, unsigned seed) { + std::mt19937 rng(seed); + std::uniform_int_distribution nnz_dist(3, 12); + std::uniform_int_distribution term_dist(0, dim - 1); + std::uniform_real_distribution val_dist(0.05F, 1.0F); + RawCsr 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; +} + +// Batch search over the corpus rows as queries; returns per-query {labels, +// scores} for a bit-exact comparison. +std::pair, std::vector> search_corpus( + Index* index, const RawCsr& queries, int k) { + std::vector labels(static_cast(queries.n) * k, + detail::INVALID_IDX); + std::vector distances(static_cast(queries.n) * k, -1.0F); + index->search(queries.n, queries.indptr.data(), queries.indices.data(), + queries.values.data(), k, distances.data(), labels.data()); + return {labels, distances}; +} + } // namespace TEST(SeismicIndexMmapIO, mapped_read_matches_the_copying_read) { @@ -1003,6 +1050,51 @@ TEST(SeismicIndexMmapIO, mapped_read_borrows_from_the_file) { sizeof(idx_t))); } +// Building from a native CSR borrowed via mmap (read_csr(kMmap)) must match +// building the same corpus fed through add(): convert -> read_csr(kMmap) -> +// build -> persist -> mmap-reload -> search is bit-exact to the add()-fed build. +// This pins the mmap-CSR build path for the regular (unquantized) index; its +// num_vectors() derives from get_vectors(), so no count fix is needed here. +TEST(SeismicIndexMmapIO, mmap_csr_build_matches_add_build) { + constexpr int kDim = 128; + const RawCsr corpus = make_raw_corpus(800, kDim, /*seed=*/11); + const RawCsr queries = make_raw_corpus(30, kDim, /*seed=*/12); + const SeismicClusterParameters params{ + .lambda = 20, .beta = 4, .alpha = 0.5F, .seed = 42}; + + // Reference: the same corpus fed through add(). + SeismicIndex added(kDim, params); + added.add(corpus.n, corpus.indptr.data(), corpus.indices.data(), + corpus.values.data()); + added.build(); + const auto expected = search_corpus(&added, queries, 10); + + // Under test: build from a native CSR of the same corpus borrowed via mmap, + // then persist and reload. + csr_test::TempCsrFiles csr("nsparse_seis_src"); + csr_test::write_interchange_csr(csr.interchange(), corpus, kDim); + csr_layout::convert(csr.interchange(), csr.native()); + + SeismicIndex mapped_build(kDim, params); + mapped_build.read_csr(csr.native().c_str(), Residency::kMmap); + ASSERT_EQ(mapped_build.num_vectors(), static_cast(corpus.n)); + mapped_build.build(); + + TempIndexFile file("nsparse_seis_mmapbuild.idx"); + write_index(&mapped_build, file.c_str()); + std::unique_ptr reloaded( + read_index(file.c_str(), IndexIoFlag::kUseMmap)); + ASSERT_NE(reloaded, nullptr); + EXPECT_EQ(reloaded->num_vectors(), static_cast(corpus.n)); + + const auto got = search_corpus(reloaded.get(), queries, 10); + EXPECT_EQ(got.first, expected.first) << "labels differ"; + ASSERT_EQ(got.second.size(), expected.second.size()); + for (size_t i = 0; i < got.second.size(); ++i) { + EXPECT_FLOAT_EQ(got.second[i], expected.second[i]) << "score at " << i; + } +} + // A stream has no file to map, so the flag alone must not send read_index down // the mapped path. TEST(SeismicIndexMmapIO, buffered_read_copies_even_when_the_flag_is_set) {