From 723ed11dbc863e0e8a0029b2a535f2b73a222b20 Mon Sep 17 00:00:00 2001 From: Zirui Song Date: Wed, 26 Aug 2026 09:53:05 +0000 Subject: [PATCH] DiskSeismic: add python tests & de-dup serialization Signed-off-by: Zirui Song --- .gitignore | 7 +- nsparse/cluster/inverted_list_clusters.cpp | 22 +++ nsparse/cluster/inverted_list_clusters.h | 8 + nsparse/disk_seismic_index.cpp | 5 +- nsparse/io/seismic_invlists_writer.h | 16 +- python_tests/test_disk_seismic_index.py | 218 +++++++++++++++++++++ tests/seismic_invlists_writer_test.cpp | 50 +++++ 7 files changed, 321 insertions(+), 5 deletions(-) create mode 100644 python_tests/test_disk_seismic_index.py diff --git a/.gitignore b/.gitignore index 3493898..d94e85b 100644 --- a/.gitignore +++ b/.gitignore @@ -47,8 +47,13 @@ venv/ # third_party third_party/ -# build +# build (and out-of-tree build directories) build/ +build-*/ + +# SWIG copies the base interface into these per-SIMD files at build time +# (configure_file COPYONLY); they are generated, not source. +nsparse/python/swignsparse_*.swig # claude .claude/settings.local.json diff --git a/nsparse/cluster/inverted_list_clusters.cpp b/nsparse/cluster/inverted_list_clusters.cpp index b2b73ba..1bb2e68 100644 --- a/nsparse/cluster/inverted_list_clusters.cpp +++ b/nsparse/cluster/inverted_list_clusters.cpp @@ -208,6 +208,11 @@ InvertedListClusters::InvertedListClusters( } auto InvertedListClusters::get_docs(idx_t idx) const -> std::span { + // A summaries-only load leaves offsets_ empty; return an empty span rather + // than indexing it out of bounds. Populated lists are unaffected. + if (offsets_.size() < static_cast(idx) + 2) { + return {}; + } return {docs_.data() + offsets_[idx], static_cast(offsets_[idx + 1] - offsets_[idx])}; } @@ -340,6 +345,23 @@ void InvertedListClusters::serialize(IOWriter* writer) const { writer->write(&n_offsets, sizeof(size_t), 1); io_align::write_padded(writer, offsets_.data(), n_offsets); + serialize_summary_store(writer); +} + +void InvertedListClusters::serialize_summaries_only(IOWriter* writer) const { + // docs_/offsets_ emitted as count-0 arrays: the membership is redundant + // with the inline forward index and unread on the mmap path. Count-0 keeps + // the byte layout identical (io/align.h), so the reader is unchanged. + size_t zero = 0; + writer->write(&zero, sizeof(size_t), 1); + io_align::write_padded(writer, nullptr, 0); + writer->write(&zero, sizeof(size_t), 1); + io_align::write_padded(writer, nullptr, 0); + + serialize_summary_store(writer); +} + +void InvertedListClusters::serialize_summary_store(IOWriter* writer) const { // Transposed (CSC) summary store. size_t n_clusters = n_clusters_; writer->write(&n_clusters, sizeof(size_t), 1); diff --git a/nsparse/cluster/inverted_list_clusters.h b/nsparse/cluster/inverted_list_clusters.h index 3281070..8839c11 100644 --- a/nsparse/cluster/inverted_list_clusters.h +++ b/nsparse/cluster/inverted_list_clusters.h @@ -40,6 +40,10 @@ class InvertedListClusters : public MmapSerializable { size_t cluster_size() const { return n_clusters_; } 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 + // regular deserialize()/mmap_deserialize(). + void serialize_summaries_only(IOWriter* writer) const; void deserialize(IOReader* reader) override; void mmap_deserialize(MmapCursor* cursor) override; @@ -54,6 +58,10 @@ class InvertedListClusters : public MmapSerializable { std::vector& out) const; private: + // Writes the transposed (CSC) summary store — the tail shared by serialize() + // and serialize_summaries_only(). + void serialize_summary_store(IOWriter* writer) const; + // Build the term-major (CSC) transpose from a per-cluster CSR summary. The // CSR summary is transient; only the transpose is retained. void build_transpose(const SparseVectors& summaries); diff --git a/nsparse/disk_seismic_index.cpp b/nsparse/disk_seismic_index.cpp index 369ef8a..4b5b9bf 100644 --- a/nsparse/disk_seismic_index.cpp +++ b/nsparse/disk_seismic_index.cpp @@ -297,7 +297,10 @@ auto DiskSeismicIndex::single_query( void DiskSeismicIndex::write_index(IOWriter* io_writer) { const uint64_t nv = num_vectors_; io_writer->write(const_cast(&nv), sizeof(uint64_t), 1); - SeismicInvertedListsWriter inv_list_writer(clustered_inverted_lists); + // 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 diff --git a/nsparse/io/seismic_invlists_writer.h b/nsparse/io/seismic_invlists_writer.h index 4acef34..9965d56 100644 --- a/nsparse/io/seismic_invlists_writer.h +++ b/nsparse/io/seismic_invlists_writer.h @@ -26,9 +26,13 @@ namespace nsparse { class SeismicInvertedListsWriter : public MmapSerializable { public: // For writing. `clustered_inverted_lists` must outlive this writer. + // `summaries_only` writes each list's doc-id membership empty (see + // InvertedListClusters::serialize_summaries_only). explicit SeismicInvertedListsWriter( - const std::vector& clustered_inverted_lists) - : borrowed_(&clustered_inverted_lists) {} + const std::vector& clustered_inverted_lists, + bool summaries_only = false) + : borrowed_(&clustered_inverted_lists), + summaries_only_(summaries_only) {} // For reading; deserialize() fills the internal store. SeismicInvertedListsWriter() = default; @@ -38,7 +42,11 @@ class SeismicInvertedListsWriter : public MmapSerializable { size_t size = lists.size(); writer->write(&size, sizeof(size), 1); for (const auto& clusters : lists) { - clusters.serialize(writer); + if (summaries_only_) { + clusters.serialize_summaries_only(writer); + } else { + clusters.serialize(writer); + } } } void deserialize(IOReader* reader) override { @@ -67,6 +75,8 @@ class SeismicInvertedListsWriter : public MmapSerializable { // when default-constructed and filled by deserialize(). const std::vector* borrowed_ = nullptr; std::vector owned_; + // Write-side only: omit the doc-id membership (reader is agnostic). + bool summaries_only_ = false; }; } // namespace nsparse diff --git a/python_tests/test_disk_seismic_index.py b/python_tests/test_disk_seismic_index.py new file mode 100644 index 0000000..85253b3 --- /dev/null +++ b/python_tests/test_disk_seismic_index.py @@ -0,0 +1,218 @@ +# 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 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. +""" + +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, +) + +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() + + +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]) diff --git a/tests/seismic_invlists_writer_test.cpp b/tests/seismic_invlists_writer_test.cpp index c01a6ce..4fbbf74 100644 --- a/tests/seismic_invlists_writer_test.cpp +++ b/tests/seismic_invlists_writer_test.cpp @@ -135,6 +135,56 @@ TEST(SeismicInvertedListsWriter, serialize_deserialize_with_summaries) { ASSERT_EQ(result[0].cluster_size(), 2); } +// summaries_only omits the doc-id membership: the stream must be smaller yet +// still round-trip through the unchanged reader with the summaries intact. +TEST(SeismicInvertedListsWriter, summaries_only_drops_membership_keeps_summaries) { + std::vector> docs = {{0, 1}, {2}}; + nsparse::InvertedListClusters clusters(docs); + auto vectors = create_float_vectors( + {{0, 1}, {0, 2}, {1, 2}}, {{1.0F, 2.0F}, {1.5F, 1.0F}, {3.0F, 2.0F}}); + clusters.summarize(&vectors, 1.0F); + std::vector original; + original.push_back(std::move(clusters)); + + nsparse::BufferedIOWriter full_writer; + nsparse::SeismicInvertedListsWriter(original).serialize(&full_writer); + nsparse::BufferedIOWriter summaries_writer; + nsparse::SeismicInvertedListsWriter(original, /*summaries_only=*/true) + .serialize(&summaries_writer); + + // Dropping docs_/offsets_ makes the summaries-only stream strictly smaller. + EXPECT_LT(summaries_writer.size(), full_writer.size()); + + // Both streams round-trip through the same (unchanged) reader. + nsparse::BufferedIOReader full_reader(full_writer.data()); + nsparse::SeismicInvertedListsWriter full_rt; + full_rt.deserialize(&full_reader); + auto full_lists = full_rt.release(); + + nsparse::BufferedIOReader summaries_reader(summaries_writer.data()); + nsparse::SeismicInvertedListsWriter summaries_rt; + summaries_rt.deserialize(&summaries_reader); + auto summaries_lists = summaries_rt.release(); + + ASSERT_EQ(full_lists.size(), 1U); + ASSERT_EQ(summaries_lists.size(), 1U); + // The cluster count (a stored scalar) survives; only membership was dropped. + EXPECT_EQ(summaries_lists[0].cluster_size(), full_lists[0].cluster_size()); + + // Summaries are byte-preserved: the same query scores both identically. + std::vector q_idx = {0, 1, 2}; + std::vector q_val = {1.0F, 1.0F, 1.0F}; + std::vector full_scores; + std::vector summaries_scores; + full_lists[0].score_summaries_transposed( + q_idx.data(), reinterpret_cast(q_val.data()), + q_idx.size(), full_scores); + summaries_lists[0].score_summaries_transposed( + q_idx.data(), reinterpret_cast(q_val.data()), + q_idx.size(), summaries_scores); + EXPECT_EQ(summaries_scores, full_scores); +} + // release() hands over what deserialize() read; it is not a way to take back a // vector handed in for writing, which the writer only borrows. TEST(SeismicInvertedListsWriter, release_moves_deserialized_data) {