diff --git a/include/svs/concurrent/dynamic_index.h b/include/svs/concurrent/dynamic_index.h index 69aee51f..96fa7c9d 100644 --- a/include/svs/concurrent/dynamic_index.h +++ b/include/svs/concurrent/dynamic_index.h @@ -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" @@ -141,6 +142,9 @@ enum class ReplaceExternalIdResult : uint8_t { NewIdExists, }; +inline constexpr size_t no_external_id = + type_traits::sentinel_v, std::less<>>.id(); + class ValidBuilder { public: ValidBuilder(const lib::SegmentedVector& status) @@ -167,6 +171,7 @@ class ValidBuilder { template class MutableVamanaIndex { friend class MultiMutableVamanaIndex; + template friend class BatchIterator; public: // Traits @@ -248,6 +253,8 @@ class MutableVamanaIndex { // (add_points takes them sequentially), so they have no relative order. std::unique_ptr compact_mutex_{ std::make_unique()}; + // 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 slot_alloc_mutex_{std::make_unique()}; @@ -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_}; @@ -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(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); } /// @@ -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 requires(std::tuple_size_v == 2) @@ -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 + requires std::invocable 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. /// @@ -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 lock_for_search() const { return std::shared_lock(*compact_mutex_); } @@ -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 void search( const Query& query, @@ -787,11 +811,9 @@ class MutableVamanaIndex { const search_parameters_type& sp, const lib::DefaultPredicate& cancel = lib::Returns(lib::Const()) ) { + // 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()}, @@ -1376,6 +1398,7 @@ class MutableVamanaIndex { re->set_recording(true); } graph_.rebuild_reverse_edges(threadpool_); + ++compaction_epoch_; } ///// Threading Interface diff --git a/include/svs/concurrent/iterator.h b/include/svs/concurrent/iterator.h index 37c0ba5c..592e5f57 100644 --- a/include/svs/concurrent/iterator.h +++ b/include/svs/concurrent/iterator.h @@ -104,14 +104,18 @@ template 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(); } @@ -226,11 +230,21 @@ template 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; } @@ -257,19 +271,20 @@ template 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([&]( const auto& graph, const auto& data, @@ -343,7 +358,8 @@ template class BatchIterator { std::vector query_; // Local buffer for the query. scratchspace_type scratchspace_; // Scratch space for search. std::vector> results_{}; // Filtered results from search. - std::unordered_set yielded_{}; // Set of yielded neighbors. + std::unordered_set 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_ = diff --git a/include/svs/concurrent/multi.h b/include/svs/concurrent/multi.h index f8f41506..f8dba7c1 100644 --- a/include/svs/concurrent/multi.h +++ b/include/svs/concurrent/multi.h @@ -134,8 +134,18 @@ template 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 contents() const { return lib::as_const_span(results_); } @@ -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 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 std::vector add_points(const Points& points, const Labels& labels, bool reuse_empty = false) { @@ -529,7 +557,9 @@ class MultiMutableVamanaIndex { for (; j < num_neighbors; ++j) { // insert default neighbor if not enough - results.set(Neighbor{}, i, j); + results.set( + type_traits::sentinel_v, compare>, i, j + ); } } } @@ -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. diff --git a/tests/svs/concurrent/concurrency.cpp b/tests/svs/concurrent/concurrency.cpp index 96a7ddf5..bc5719e8 100644 --- a/tests/svs/concurrent/concurrency.cpp +++ b/tests/svs/concurrent/concurrency.cpp @@ -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) {