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
7 changes: 6 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
22 changes: 22 additions & 0 deletions nsparse/cluster/inverted_list_clusters.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,11 @@ InvertedListClusters::InvertedListClusters(
}

auto InvertedListClusters::get_docs(idx_t idx) const -> std::span<const idx_t> {
// 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<size_t>(idx) + 2) {
return {};
}
return {docs_.data() + offsets_[idx],
static_cast<size_t>(offsets_[idx + 1] - offsets_[idx])};
}
Expand Down Expand Up @@ -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<idx_t>(writer, nullptr, 0);
writer->write(&zero, sizeof(size_t), 1);
io_align::write_padded<idx_t>(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);
Expand Down
8 changes: 8 additions & 0 deletions nsparse/cluster/inverted_list_clusters.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -54,6 +58,10 @@ class InvertedListClusters : public MmapSerializable {
std::vector<float>& 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);
Expand Down
5 changes: 4 additions & 1 deletion nsparse/disk_seismic_index.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<uint64_t*>(&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
Expand Down
16 changes: 13 additions & 3 deletions nsparse/io/seismic_invlists_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -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<InvertedListClusters>& clustered_inverted_lists)
: borrowed_(&clustered_inverted_lists) {}
const std::vector<InvertedListClusters>& clustered_inverted_lists,
bool summaries_only = false)
: borrowed_(&clustered_inverted_lists),
summaries_only_(summaries_only) {}

// For reading; deserialize() fills the internal store.
SeismicInvertedListsWriter() = default;
Expand All @@ -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 {
Expand Down Expand Up @@ -67,6 +75,8 @@ class SeismicInvertedListsWriter : public MmapSerializable {
// when default-constructed and filled by deserialize().
const std::vector<InvertedListClusters>* borrowed_ = nullptr;
std::vector<InvertedListClusters> owned_;
// Write-side only: omit the doc-id membership (reader is agnostic).
bool summaries_only_ = false;
};
} // namespace nsparse

Expand Down
218 changes: 218 additions & 0 deletions python_tests/test_disk_seismic_index.py
Original file line number Diff line number Diff line change
@@ -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])
50 changes: 50 additions & 0 deletions tests/seismic_invlists_writer_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<std::vector<nsparse::idx_t>> 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<nsparse::InvertedListClusters> 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<nsparse::term_t> q_idx = {0, 1, 2};
std::vector<float> q_val = {1.0F, 1.0F, 1.0F};
std::vector<float> full_scores;
std::vector<float> summaries_scores;
full_lists[0].score_summaries_transposed(
q_idx.data(), reinterpret_cast<const uint8_t*>(q_val.data()),
q_idx.size(), full_scores);
summaries_lists[0].score_summaries_transposed(
q_idx.data(), reinterpret_cast<const uint8_t*>(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) {
Expand Down
Loading