Skip to content
Open
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
53 changes: 38 additions & 15 deletions include/svs/concurrent/dynamic_index.h
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
#include "svs/core/recall.h"
#include "svs/index/vamana/index.h"
#include "svs/lib/boundscheck.h"
#include "svs/lib/neighbor.h"
#include "svs/lib/preprocessor.h"
#include "svs/lib/scopeguard.h"
#include "svs/lib/segmented_vector.h"
Expand Down Expand Up @@ -141,6 +142,9 @@ enum class ReplaceExternalIdResult : uint8_t {
NewIdExists,
};

inline constexpr size_t no_external_id =
type_traits::sentinel_v<Neighbor<size_t>, std::less<>>.id();

class ValidBuilder {
public:
ValidBuilder(const lib::SegmentedVector<SlotMetadata>& status)
Expand All @@ -167,6 +171,7 @@ class ValidBuilder {
template <graphs::MemoryGraph Graph, typename Data, typename Dist>
class MutableVamanaIndex {
friend class MultiMutableVamanaIndex<Graph, Data, Dist>;
template <typename, typename> friend class BatchIterator;

public:
// Traits
Expand Down Expand Up @@ -248,6 +253,8 @@ class MutableVamanaIndex {
// (add_points takes them sequentially), so they have no relative order.
std::unique_ptr<std::shared_mutex> compact_mutex_{
std::make_unique<std::shared_mutex>()};
// Iterators check this under compact_mutex_ before reusing cached internal IDs.
size_t compaction_epoch_ = 0;
// Writer-only mutex serializing slot allocation in add_points
std::unique_ptr<std::mutex> slot_alloc_mutex_{std::make_unique<std::mutex>()};

Expand Down Expand Up @@ -571,7 +578,8 @@ class MutableVamanaIndex {
///
/// @param i The internal ID to translate to an external ID.
///
/// Requires that mapping for `i` exists. Otherwise, all bets are off.
/// Returns `no_external_id` if `i` has no mapping (a concurrent consolidation erased
/// it).
///
size_t translate_internal_id(Idx i) const {
std::shared_lock lock{*translator_mutex_};
Expand All @@ -581,9 +589,8 @@ class MutableVamanaIndex {
/// @copydoc translate_internal_id
/// Requires the caller to hold `lock_for_translation()`.
size_t unsafe_translate_internal_id(Idx i) const {
// Use get_external_or to handle concurrent consolidate erasing entries.
// If the entry was erased, return the internal ID as-is (stale result).
return translator_.get_external_or(i, static_cast<size_t>(i));
// Not `i` itself: an internal ID is no label, and a caller would report it as one.
return translator_.get_external_or(i, no_external_id);
}

///
Expand Down Expand Up @@ -641,8 +648,9 @@ class MutableVamanaIndex {
/// (1) This is definitely not safe to call multiple times on the same array for obvious
/// reasons.
///
/// (2) All entries in `ids` should have valid translations. Otherwise, this function's
/// behavior is undefined.
/// (2) An entry without a translation becomes `no_external_id`; callers skip it.
///
/// (3) Exclude compaction from the production of these internal IDs through this call.
///
template <class Dims, class Base>
requires(std::tuple_size_v<Dims> == 2)
Expand Down Expand Up @@ -697,6 +705,23 @@ class MutableVamanaIndex {
return data_.get_datum(translator_.get_internal(e));
}

///
/// @brief Call `f` with the raw data of external id `e` while it is locked.
///
/// @returns Whether `e` was live
///
template <typename F>
requires std::invocable<F&, typename Data::const_value_type> bool
on_datum(size_t e, F&& f) const {
std::shared_lock compact_lock{*compact_mutex_};
std::shared_lock lock{*translator_mutex_};
if (!unsafe_has_id(e)) {
return false;
}
f(data_.get_datum(translator_.get_internal(e)));
return true;
}

///
/// @brief Return the dimensionality of the stored dataset.
///
Expand All @@ -714,10 +739,8 @@ class MutableVamanaIndex {
/// greedy traversal — mirrors the shared lock taken by search(). Growth by
/// add_points needs no lock (grow-stable SegmentedVector storage).
///
/// Acquire this only around graph traversal, and release it before
/// acquiring lock_for_translation(): the two must never be held nested in
/// the compact->translator order reversed, which would invert the global
/// lock order (compact -> translator) and deadlock against compact.
/// Hold this through internal-to-external ID translation. Acquire the translation
/// lock inside this guard, preserving the compact -> translator lock order.
[[nodiscard]] std::shared_lock<std::shared_mutex> lock_for_search() const {
return std::shared_lock<std::shared_mutex>(*compact_mutex_);
}
Expand Down Expand Up @@ -760,7 +783,8 @@ class MutableVamanaIndex {
};
}

// Single Search
// The scratch buffer contains internal IDs. Consuming them after this call requires
// external exclusion of compaction; use the result-view overload for external IDs.
template <typename Query>
void search(
const Query& query,
Expand All @@ -787,11 +811,9 @@ class MutableVamanaIndex {
const search_parameters_type& sp,
const lib::DefaultPredicate& cancel = lib::Returns(lib::Const<false>())
) {
// Internal IDs must retain their meaning until translation has completed.
std::shared_lock compact_lock{*compact_mutex_};
{
// compact_mutex_ shared: blocks compact()'s segment-freeing shrink
// during the traversal. Released before translate_to_external() takes
// translator_mutex_ to keep the compact->translator lock order.
std::shared_lock compact_lock{*compact_mutex_};
threads::parallel_for(
threadpool_,
threads::StaticPartition{queries.size()},
Expand Down Expand Up @@ -1376,6 +1398,7 @@ class MutableVamanaIndex {
re->set_recording(true);
}
graph_.rebuild_reverse_edges(threadpool_);
++compaction_epoch_;
}

///// Threading Interface
Expand Down
48 changes: 32 additions & 16 deletions include/svs/concurrent/iterator.h
Original file line number Diff line number Diff line change
Expand Up @@ -104,14 +104,18 @@ template <typename Index, typename QueryType> class BatchIterator {
results_.clear();
const auto& buffer = scratchspace_.buffer;
for (size_t i = 0, imax = buffer.size(); i < imax; ++i) {
auto neighbor = buffer[i];
auto neighbor = adapt(buffer[i]);
// Erased by a concurrent consolidation.
if (neighbor.id() == no_external_id) {
continue;
}
auto result = yielded_.insert(neighbor.id());
if (result.second /* inserted */) {
// Rollback insertion into the yielded set if push_back throws.
auto guard = lib::make_dismissable_scope_guard([&]() noexcept {
yielded_.erase(result.first);
});
results_.push_back(adapt(neighbor));
results_.push_back(neighbor);
guard.dismiss();
}

Expand Down Expand Up @@ -226,11 +230,21 @@ template <typename Index, typename QueryType> class BatchIterator {
/// @brief Returns whether iterator can find more neighbors or not for the given query.
///
/// The iterator is considered done when all the available nodes have been yielded or
/// when the search can not find any more neighbors. The transition from not done to
/// done will be triggered by a call to ``next()``. The contents of ``batch_number()``
/// and ``parameters_for_current_iteration()`` will then remain unchanged by subsequent
/// invocations of ``next()``.
bool done() const { return (yielded_.size() == parent_->size() || is_exhausted_); }
/// when the search can not find any more neighbors. Deleted IDs in the yielded set
/// do not count as coverage of the current live IDs.
bool done() const {
if (is_exhausted_) {
return true;
}
if (yielded_.size() < parent_->size()) {
return false;
}
bool all_yielded = true;
parent_->on_ids([&](size_t id) {
all_yielded = all_yielded && yielded_.contains(id);
});
return all_yielded;
}

/// @brief Forces the next iteration to restart the search from scratch.
void restart_next_search() { restart_search_ = true; }
Expand All @@ -257,19 +271,20 @@ template <typename Index, typename QueryType> class BatchIterator {
return;
}

// Keep internal IDs stable through traversal, reranking, and translation.
[[maybe_unused]] auto search_guard = parent_->lock_for_search();
if (compaction_epoch_ != parent_->compaction_epoch_) {
// Reranking touches cached IDs before the restart initializer runs.
// Preserve the grown window so the restart can reach unreturned neighbors.
scratchspace_.buffer.clear();
restart_search_ = true;
compaction_epoch_ = parent_->compaction_epoch_;
}
increment_buffer(batch_size);

bool restart_search_copy = std::exchange(restart_search_, true);

// Hold the search lock (compact_mutex_ shared for a dynamic index; a
// no-op for a static one) so a concurrent compact() cannot free the
// segments under us. add_points growth is lock-free (grow-stable
// storage). The guard is released at the end of this scope — before
// acquiring the translation lock. The two locks must never be held
// nested in the translator-before-compact order: that would invert the
// global lock order (compact -> translator) and deadlock against compact.
{
[[maybe_unused]] auto search_guard = parent_->lock_for_search();
parent_->experimental_escape_hatch([&]<std::integral I>(
const auto& graph,
const auto& data,
Expand Down Expand Up @@ -343,7 +358,8 @@ template <typename Index, typename QueryType> class BatchIterator {
std::vector<QueryType> query_; // Local buffer for the query.
scratchspace_type scratchspace_; // Scratch space for search.
std::vector<Neighbor<external_id_type>> results_{}; // Filtered results from search.
std::unordered_set<internal_id_type> yielded_{}; // Set of yielded neighbors.
std::unordered_set<external_id_type> yielded_{}; // Set of yielded neighbors.
size_t compaction_epoch_ = 0; // Epoch of scratchspace_ IDs.
size_t iteration_ = 0; // Current iteration number.
bool restart_search_ = true; // Whether the next search should restart from scratch.
size_t extra_search_buffer_capacity_ =
Expand Down
40 changes: 34 additions & 6 deletions include/svs/concurrent/multi.h
Original file line number Diff line number Diff line change
Expand Up @@ -134,8 +134,18 @@ template <typename Index, typename QueryType> class MultiBatchIterator {
size_t size() const { return results_.size(); }

bool done() const {
return (batch_iterator_.done() && extra_results_.empty()) ||
(returned_.size() == index_.labelcount());
if (batch_iterator_.done() && extra_results_.empty()) {
return true;
}
if (returned_.size() < index_.labelcount()) {
return false;
}
// Deletions can make the counts equal while live labels remain unreturned.
bool all_returned = true;
index_.on_ids([&](label_type label) {
all_returned = all_returned && returned_.contains(label);
});
return all_returned;
}

std::span<const value_type> contents() const { return lib::as_const_span(results_); }
Expand Down Expand Up @@ -400,6 +410,24 @@ class MultiMutableVamanaIndex {
return best;
}

///
/// @brief Call `f` with the raw data of every vector of `label` while it is locked.
///
/// @returns The number of vectors `f` was called for.
///
template <typename F> size_t on_data(label_type label, F&& f) const {
std::shared_lock l2e_lock{*l2e_mutex_};
auto it = label_to_external_.find(label);
if (it == label_to_external_.end()) {
return 0;
}
size_t visited = 0;
for (auto external_id : it->second) {
visited += index_->on_datum(external_id, f);
}
return visited;
}

template <data::ImmutableMemoryDataset Points, typename Labels>
std::vector<external_id_type>
add_points(const Points& points, const Labels& labels, bool reuse_empty = false) {
Expand Down Expand Up @@ -529,7 +557,9 @@ class MultiMutableVamanaIndex {

for (; j < num_neighbors; ++j) {
// insert default neighbor if not enough
results.set(Neighbor<label_type>{}, i, j);
results.set(
type_traits::sentinel_v<Neighbor<label_type>, compare>, i, j
);
}
}
}
Expand Down Expand Up @@ -629,9 +659,7 @@ class MultiMutableVamanaIndex {
// translate internal id -> external id -> label
label_type translate_internal_id(Idx i) const {
auto external_id = index_->translate_internal_id(i);
// Mirrors the parent's get_external_or: fall back to the external id when the
// label mapping is already gone.
return find_label(external_id).value_or(external_id);
return find_label(external_id).value_or(no_external_id);
}

/// @brief The label owning `external_id`, or ``std::nullopt`` if it has none.
Expand Down
4 changes: 2 additions & 2 deletions tests/svs/concurrent/concurrency.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -308,8 +308,8 @@ CATCH_TEST_CASE(
// Every ID a search hands back either still maps to a live external
// ID -- in which case the mapping must round-trip exactly -- or it
// names a slot retired by a concurrent deleter, which is expected
// and unobservable from here (``translate_internal_id`` degrades to
// returning the internal ID when the entry has been erased).
// and unobservable from here (``translate_internal_id`` returns
// ``no_external_id`` when the entry has been erased).
auto external = index->translate_internal_id(internal);
if (index->has_id(external) &&
index->translate_external_id(external) != internal) {
Expand Down
Loading