Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions nsparse/disk_seismic_index_base.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand Down
80 changes: 80 additions & 0 deletions tests/csr_interchange_test_util.h
Original file line number Diff line number Diff line change
@@ -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 <array>
#include <cstdint>
#include <filesystem>
#include <fstream>
#include <string>
#include <system_error>
#include <vector>

// 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 <class Corpus>
void write_interchange_csr(const std::string& path, const Corpus& c,
int num_cols) {
std::ofstream out(path, std::ios::binary);
const std::array<int64_t, 3> header = {
static_cast<int64_t>(c.n), static_cast<int64_t>(num_cols),
static_cast<int64_t>(c.indices.size())};
out.write(reinterpret_cast<const char*>(header.data()),
header.size() * sizeof(int64_t));
const std::vector<int64_t> indptr64(c.indptr.begin(), c.indptr.end());
out.write(reinterpret_cast<const char*>(indptr64.data()),
static_cast<std::streamsize>(indptr64.size() * sizeof(int64_t)));
const std::vector<int32_t> indices32(c.indices.begin(), c.indices.end());
out.write(reinterpret_cast<const char*>(indices32.data()),
static_cast<std::streamsize>(indices32.size() * sizeof(int32_t)));
out.write(reinterpret_cast<const char*>(c.values.data()),
static_cast<std::streamsize>(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
43 changes: 43 additions & 0 deletions tests/disk_seismic_index_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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<Index> added_mapped(
read_index(added_file.c_str(), IndexIoFlag::kUseMmap));
ASSERT_NE(added_mapped, nullptr);
const ScoreIds fresh = search_all(*added_mapped, queries, 10, &params);

// 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<size_t>(corpus.n));
mapped_build.build();
TempIndexFile built_file("nsparse_disk_seismic_mmapbuild.idx");
write_index(&mapped_build, built_file.c_str());
std::unique_ptr<Index> 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<size_t>(corpus.n));
expect_same_results(search_all(*built_mapped, queries, 10, &params), 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.
Expand Down
6 changes: 6 additions & 0 deletions tests/disk_seismic_test_util.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down
92 changes: 92 additions & 0 deletions tests/seismic_index_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
#include <map>
#include <memory>
#include <random>
#include <set>
#include <string>
#include <unordered_set>
#include <vector>
Expand All @@ -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 {
Expand Down Expand Up @@ -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<idx_t> indptr;
std::vector<term_t> indices;
std::vector<float> values;
idx_t n = 0;
};

RawCsr make_raw_corpus(idx_t rows, int dim, unsigned seed) {
std::mt19937 rng(seed);
std::uniform_int_distribution<int> nnz_dist(3, 12);
std::uniform_int_distribution<int> term_dist(0, dim - 1);
std::uniform_real_distribution<float> 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<int> terms;
const int nnz = nnz_dist(rng);
while (static_cast<int>(terms.size()) < nnz) {
terms.insert(term_dist(rng));
}
for (const int t : terms) {
c.indices.push_back(static_cast<term_t>(t));
c.values.push_back(val_dist(rng));
}
c.indptr.push_back(static_cast<idx_t>(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<idx_t>, std::vector<float>> search_corpus(
Index* index, const RawCsr& queries, int k) {
std::vector<idx_t> labels(static_cast<size_t>(queries.n) * k,
detail::INVALID_IDX);
std::vector<float> distances(static_cast<size_t>(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) {
Expand Down Expand Up @@ -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<size_t>(corpus.n));
mapped_build.build();

TempIndexFile file("nsparse_seis_mmapbuild.idx");
write_index(&mapped_build, file.c_str());
std::unique_ptr<Index> reloaded(
read_index(file.c_str(), IndexIoFlag::kUseMmap));
ASSERT_NE(reloaded, nullptr);
EXPECT_EQ(reloaded->num_vectors(), static_cast<size_t>(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) {
Expand Down
Loading