From 02300e05911702edd305bbb5b8aab6caa52291d3 Mon Sep 17 00:00:00 2001 From: Kevin Boyd Date: Tue, 6 Oct 2026 12:14:36 -0400 Subject: [PATCH 1/4] Run SubstructLibrary queries concurrently within a GPU-memory budget Instead of serializing every call, let queries run concurrently. Each finalize() works out how many queries fit below 85% of each GPU's memory, counting the batches it is about to replace as freed, and gives every GPU that many query slots, each a search workspace plus its match flags. A query takes one slot per GPU, and queries with recursive SMARTS also take one of a smaller number of slots for their scratch memory. maxConcurrentQueries() reports the limit. Queries share the library lock; additions and finalize() take it exclusively, passing a gate first so queries that keep arriving cannot starve them. If a finalize() fails after resizing the slots, the previous slots are rebuilt. The memory estimates live next to the search they describe, which also now resolves the default executors per runner in one place. --- src/substruct/substruct_library.cpp | 241 ++++++++++++++++++++++++---- src/substruct/substruct_library.h | 6 +- src/substruct/substruct_search.cu | 55 +++++-- src/substruct/substruct_search.h | 10 ++ tests/test_substruct_library.cu | 90 ++++++++++- 5 files changed, 357 insertions(+), 45 deletions(-) diff --git a/src/substruct/substruct_library.cpp b/src/substruct/substruct_library.cpp index 7ead077e..9d77938c 100644 --- a/src/substruct/substruct_library.cpp +++ b/src/substruct/substruct_library.cpp @@ -9,9 +9,12 @@ #include #include +#include #include #include #include +#include +#include #include #include #include @@ -63,6 +66,65 @@ struct GpuTargets { std::vector molecules; std::vector ids; std::shared_ptr searchTargets; + + //! Approximate device memory held by the uploaded batch. + [[nodiscard]] std::size_t deviceBytes() const { + return host->batchAtomStarts.size() * sizeof(int) + host->atomDataPacked.size() * sizeof(AtomDataPacked) + + host->bondTypeCounts.size() * sizeof(BondTypeCounts) + + host->targetAtomBonds.size() * sizeof(TargetAtomBonds); + } +}; + +//! A search workspace and the per-target match flags its queries fill. +struct QuerySlot { + std::shared_ptr workspace; + std::vector matchFlags; +}; + +//! One GPU's query slots; a query waits for a free slot. +class SlotPool { + public: + explicit SlotPool(std::vector> slots) : slots_(std::move(slots)) { + for (auto& slot : slots_) { + free_.push_back(slot.get()); + } + } + + [[nodiscard]] QuerySlot* take() { + std::unique_lock lock(mutex_); + available_.wait(lock, [this] { return !free_.empty(); }); + QuerySlot* slot = free_.back(); + free_.pop_back(); + return slot; + } + + void give(QuerySlot* slot) { + { + const std::lock_guard lock(mutex_); + free_.push_back(slot); + } + available_.notify_one(); + } + + private: + std::vector> slots_; + std::mutex mutex_; + std::condition_variable available_; + std::vector free_; +}; + +//! Holds a slot taken from a pool for the duration of one GPU search. +class TakenSlot { + public: + explicit TakenSlot(SlotPool& pool) : pool_(pool), slot_(pool.take()) {} + ~TakenSlot() { pool_.give(slot_); } + TakenSlot(const TakenSlot&) = delete; + TakenSlot& operator=(const TakenSlot&) = delete; + [[nodiscard]] QuerySlot& operator*() const { return *slot_; } + + private: + SlotPool& pool_; + QuerySlot* slot_; }; } // namespace @@ -90,21 +152,18 @@ class SubstructLibrary::Impl { gpus_.size(), [&](std::size_t index) { return gpus_[index] != nullptr; }, [&](std::size_t index) { gpus_[index].reset(); }); - releaseOnDevices( - workspaces_.size(), - [&](std::size_t index) { return workspaces_[index] != nullptr; }, - [&](std::size_t index) { workspaces_[index].reset(); }); + destroySlotPools(); } unsigned int addMol(const RDKit::ROMol& molecule) { - const std::lock_guard lock(mutex_); + const auto lock = lockForWriting(); requireIdsAvailable(1); molecules_.push_back(copyWithRingInfo(molecule)); return static_cast(molecules_.size() - 1); } std::vector addMols(const std::vector& molecules) { - const std::lock_guard lock(mutex_); + const auto lock = lockForWriting(); requireIdsAvailable(molecules.size()); const std::size_t firstId = molecules_.size(); molecules_.resize(firstId + molecules.size()); @@ -135,7 +194,7 @@ class SubstructLibrary::Impl { } void finalize(cudaStream_t stream) { - const std::lock_guard lock(mutex_); + const auto lock = lockForWriting(); useCurrentDeviceIfNoneConfigured(); if (deviceIds_.size() > 1 && stream != nullptr) { throw std::invalid_argument("A single external CUDA stream cannot be used with a multi-GPU substructure library"); @@ -208,16 +267,29 @@ class SubstructLibrary::Impl { exceptions.store(std::current_exception()); } } + bool slotsRebuilt = false; try { exceptions.rethrow(); - createWorkspaces(); rdkitIds_.reserve(rdkitIds_.size() + newRdkitIds.size()); + // Query slots are sized against memory that already holds the new batches. The batches they replace are + // freed below, so their memory counts as available. + std::vector freedBytes(deviceIds_.size(), 0); + for (std::size_t gpu = 0; gpu < deviceIds_.size(); ++gpu) { + if (replacements[gpu] != nullptr && gpus_[gpu] != nullptr) { + freedBytes[gpu] = gpus_[gpu]->deviceBytes(); + } + } + slotsRebuilt = true; + createSlotPools(freedBytes); } catch (...) { static_cast(cudaGetLastError()); releaseOnDevices( replacements.size(), [&](std::size_t index) { return replacements[index] != nullptr; }, [&](std::size_t index) { replacements[index].reset(); }); + if (slotsRebuilt && finalized_) { + restoreSlotPools(); + } throw; } @@ -231,19 +303,24 @@ class SubstructLibrary::Impl { } [[nodiscard]] std::size_t size() const { - const std::lock_guard lock(mutex_); + const auto lock = lockForQuery(); return numFinalized_; } [[nodiscard]] std::size_t pendingSize() const { - const std::lock_guard lock(mutex_); + const auto lock = lockForQuery(); return molecules_.size() - numFinalized_; } + [[nodiscard]] std::size_t maxConcurrentQueries() const { + const auto lock = lockForQuery(); + return maxConcurrentQueries_; + } + [[nodiscard]] std::vector getMatches(const RDKit::ROMol& query, int maxResults, cudaStream_t stream) const { - const std::lock_guard lock(mutex_); + const auto lock = lockForQuery(); requireFinalized(); if (maxResults == 0) { return {}; @@ -252,13 +329,13 @@ class SubstructLibrary::Impl { } [[nodiscard]] std::size_t countMatches(const RDKit::ROMol& query, cudaStream_t stream) const { - const std::lock_guard lock(mutex_); + const auto lock = lockForQuery(); requireFinalized(); return matches(query, stream, -1).size(); } [[nodiscard]] bool hasMatch(const RDKit::ROMol& query, cudaStream_t stream) const { - const std::lock_guard lock(mutex_); + const auto lock = lockForQuery(); requireFinalized(); return !matches(query, stream, 1).empty(); } @@ -274,6 +351,21 @@ class SubstructLibrary::Impl { if (!finalized_) { throw std::logic_error("Substructure library must be finalized before querying"); } + if (slotPools_.empty()) { + throw std::runtime_error("Substructure library has no query slots after a failed finalize(); retry finalize()"); + } + } + + // Writers hold writerGate_ while waiting for exclusive access and queries pass through it first, so a waiting + // writer is not starved by queries that keep arriving. + [[nodiscard]] std::unique_lock lockForWriting() { + const std::lock_guard gate(writerGate_); + return std::unique_lock(mutex_); + } + + [[nodiscard]] std::shared_lock lockForQuery() const { + { const std::lock_guard gate(writerGate_); } + return std::shared_lock(mutex_); } [[nodiscard]] int threads() const { @@ -334,17 +426,80 @@ class SubstructLibrary::Impl { return next; } - // Search state for each GPU, created on the first finalize() and reused by every query. - void createWorkspaces() { - if (workspaces_.size() == deviceIds_.size()) { - return; + // Give each GPU as many query slots as fit below 85% of its memory, counting freedBytes[gpu] as free. Queries + // with recursive SMARTS also need scratch memory, so fewer of them may run at once. + void createSlotPools(const std::vector& freedBytes) { + destroySlotPools(); + constexpr std::size_t memoryPercent = 85; + std::size_t slots = std::numeric_limits::max(); + std::vector availableBytes(deviceIds_.size()); + std::vector slotBytes(deviceIds_.size()); + std::vector recursiveBytes(deviceIds_.size()); + for (std::size_t gpu = 0; gpu < deviceIds_.size(); ++gpu) { + const WithDevice device(deviceIds_[gpu]); + std::size_t freeBytes = 0; + std::size_t totalBytes = 0; + cudaCheckError(cudaMemGetInfo(&freeBytes, &totalBytes)); + const std::size_t usedBytes = totalBytes - std::min(totalBytes, freeBytes + freedBytes[gpu]); + const std::size_t budgetBytes = totalBytes * memoryPercent / 100; + availableBytes[gpu] = budgetBytes > usedBytes ? budgetBytes - usedBytes : 0; + slotBytes[gpu] = estimateSubstructSearchWorkspaceBytes(gpuConfig(gpu)); + recursiveBytes[gpu] = estimateRecursiveScratchBytes(gpuConfig(gpu)); + if (availableBytes[gpu] < slotBytes[gpu] + recursiveBytes[gpu]) { + throw std::runtime_error("Substructure library cannot fit one query below 85% of GPU memory"); + } + slots = std::min(slots, (availableBytes[gpu] - recursiveBytes[gpu]) / slotBytes[gpu]); } - std::vector> created(deviceIds_.size()); + // Each query searches every GPU from its own host thread. + const std::size_t hostThreads = static_cast(omp_get_max_threads()) / deviceIds_.size(); + slots = std::max(1, std::min(slots, hostThreads)); + std::size_t recursiveSlots = slots; for (std::size_t gpu = 0; gpu < deviceIds_.size(); ++gpu) { - created[gpu] = makeSubstructSearchWorkspace(deviceIds_[gpu]); + recursiveSlots = std::min(recursiveSlots, (availableBytes[gpu] - slots * slotBytes[gpu]) / recursiveBytes[gpu]); + } + + std::vector> pools(deviceIds_.size()); + try { + for (std::size_t gpu = 0; gpu < deviceIds_.size(); ++gpu) { + const WithDevice device(deviceIds_[gpu]); + std::vector> gpuSlots(slots); + for (auto& slot : gpuSlots) { + slot = std::make_unique(); + slot->workspace = makeSubstructSearchWorkspace(deviceIds_[gpu]); + } + pools[gpu] = std::make_unique(std::move(gpuSlots)); + } + } catch (...) { + releaseOnDevices( + pools.size(), + [&](std::size_t gpu) { return pools[gpu] != nullptr; }, + [&](std::size_t gpu) { pools[gpu].reset(); }); + throw; } - workspaces_ = std::move(created); - matchFlags_.resize(deviceIds_.size()); + slotPools_ = std::move(pools); + recursiveQuerySlots_ = std::make_unique>( + static_cast(std::max(1, recursiveSlots))); + maxConcurrentQueries_ = slots; + } + + // After a failed finalize(), give the previously finalized molecules their query slots back. If that fails too, + // queries are refused until finalize() succeeds. + void restoreSlotPools() noexcept { + try { + createSlotPools(std::vector(deviceIds_.size(), 0)); + } catch (...) { + destroySlotPools(); + } + } + + void destroySlotPools() noexcept { + releaseOnDevices( + slotPools_.size(), + [&](std::size_t gpu) { return slotPools_[gpu] != nullptr; }, + [&](std::size_t gpu) { slotPools_[gpu].reset(); }); + slotPools_.clear(); + recursiveQuerySlots_.reset(); + maxConcurrentQueries_ = 0; } [[nodiscard]] SubstructSearchConfig gpuConfig(std::size_t gpu) const { @@ -367,6 +522,19 @@ class SubstructLibrary::Impl { if (deviceIds_.size() > 1 && stream != nullptr) { throw std::invalid_argument("A single external CUDA stream cannot be used with a multi-GPU substructure library"); } + // Recursive SMARTS hold extra scratch memory on every GPU while they run. + struct RecursiveSlot { + std::counting_semaphore<>* slots = nullptr; + ~RecursiveSlot() { + if (slots != nullptr) { + slots->release(); + } + } + } recursiveSlot; + if (hasRecursiveSmarts(&query)) { + recursiveQuerySlots_->acquire(); + recursiveSlot.slots = recursiveQuerySlots_.get(); + } std::vector> perGpu(deviceIds_.size()); detail::OpenMPExceptionRegistry exceptions; #pragma omp parallel for num_threads(static_cast(deviceIds_.size())) schedule(static) @@ -417,9 +585,10 @@ class SubstructLibrary::Impl { cudaStream_t stream) const { ScopedNvtxRange searchRange("SubstructLibrary GPU search"); const GpuTargets& targets = *gpus_[gpu]; - std::vector& flags = matchFlags_[gpu]; - const SubstructSearchConfig config = gpuConfig(gpu); - hasSubstructMatch(*targets.searchTargets, query, flags, config.algorithm, stream, config, workspaces_[gpu].get()); + const TakenSlot slot(*slotPools_[gpu]); + std::vector& flags = (*slot).matchFlags; + const SubstructSearchConfig config = gpuConfig(gpu); + hasSubstructMatch(*targets.searchTargets, query, flags, config.algorithm, stream, config, (*slot).workspace.get()); std::vector result; for (std::size_t target = 0; target < flags.size(); ++target) { if (flags[target] != 0) { @@ -451,22 +620,24 @@ class SubstructLibrary::Impl { } } - SubstructSearchConfig config_; - std::vector deviceIds_; - // Held by every call, so queries, additions, and finalize() run one at a time. - mutable std::mutex mutex_; + SubstructSearchConfig config_; + std::vector deviceIds_; + // Queries share mutex_; additions and finalize() hold it exclusively. + mutable std::shared_mutex mutex_; + mutable std::mutex writerGate_; //! Every added molecule, indexed by ID; the first numFinalized_ are searchable. std::vector> molecules_; std::size_t numFinalized_ = 0; bool finalized_ = false; - std::vector> gpus_; + std::vector> gpus_; //! Finalized molecules the GPU format cannot represent, ascending; matched with RDKit. - std::vector rdkitIds_; - std::vector> workspaces_; - //! Per-GPU match flags, reused across queries to avoid reallocating them. - mutable std::vector> matchFlags_; + std::vector rdkitIds_; + mutable std::vector> slotPools_; + //! Bounds how many queries with recursive SMARTS hold scratch memory at once. + std::unique_ptr> recursiveQuerySlots_; + std::size_t maxConcurrentQueries_ = 0; }; SubstructLibrary::SubstructLibrary(SubstructSearchConfig config) : impl_(std::make_unique(std::move(config))) {} @@ -493,6 +664,10 @@ std::size_t SubstructLibrary::pendingSize() const { return impl_->pendingSize(); } +std::size_t SubstructLibrary::maxConcurrentQueries() const { + return impl_->maxConcurrentQueries(); +} + std::vector SubstructLibrary::getMatches(const RDKit::ROMol& query, int maxResults, cudaStream_t stream) const { diff --git a/src/substruct/substruct_library.h b/src/substruct/substruct_library.h index 58013596..67309d7c 100644 --- a/src/substruct/substruct_library.h +++ b/src/substruct/substruct_library.h @@ -25,7 +25,8 @@ namespace nvMolKit { * configured, molecules are split across those GPUs and results are merged in * ID order. Molecules the GPU format cannot represent are matched with RDKit. * - * Calls are serialized: queries, additions, and finalize() run one at a time. + * Queries may run concurrently, as many as fit in GPU memory; additions and + * finalize() wait for running queries and hold off new ones. */ class SubstructLibrary { public: @@ -59,6 +60,9 @@ class SubstructLibrary { /** Number of molecules added since the last finalize(). */ [[nodiscard]] std::size_t pendingSize() const; + /** How many queries can run at once, set by finalize() from the GPU memory left after the molecules. */ + [[nodiscard]] std::size_t maxConcurrentQueries() const; + /** IDs of matching molecules, ascending. maxResults limits the count; -1 returns all, 0 returns none. */ [[nodiscard]] std::vector getMatches(const RDKit::ROMol& query, int maxResults = -1, diff --git a/src/substruct/substruct_search.cu b/src/substruct/substruct_search.cu index 49b81d4b..e8bd50c4 100644 --- a/src/substruct/substruct_search.cu +++ b/src/substruct/substruct_search.cu @@ -635,6 +635,18 @@ void runGpuCoordinator(int deviceId, namespace { +/// Executors each runner thread drives: 3 for a lone runner, 2 otherwise, unless configured. +int resolveExecutorsPerRunner(int configured, int numRunners) { + if (configured == -1) { + return numRunners == 1 ? 3 : 2; + } + if (configured < 1 || configured > kMaxExecutorsPerRunner) { + throw std::invalid_argument("executorsPerRunner must be -1 (auto) or between 1 and " + + std::to_string(kMaxExecutorsPerRunner)); + } + return configured; +} + void runPipelinedSubstructSearch(const std::vector& targets, const MoleculesHost& queriesHost, const MoleculesDevice& queriesDevice, @@ -686,15 +698,7 @@ void runPipelinedSubstructSearch(const std::vector& targets return; } - int executorsPerRunner; - if (config.executorsPerRunner == -1) { - executorsPerRunner = (numRunners == 1) ? 3 : 2; - } else if (config.executorsPerRunner < 1 || config.executorsPerRunner > kMaxExecutorsPerRunner) { - throw std::invalid_argument("executorsPerRunner must be -1 (auto) or between 1 and " + - std::to_string(kMaxExecutorsPerRunner)); - } else { - executorsPerRunner = config.executorsPerRunner; - } + const int executorsPerRunner = resolveExecutorsPerRunner(config.executorsPerRunner, numRunners); std::vector workersPerGpu(numGpus, numRunners / numGpus); for (int i = 0; i < numRunners % numGpus; ++i) { @@ -1410,4 +1414,37 @@ std::shared_ptr makeSubstructSearchWorkspace(int devic return std::make_shared(deviceId); } +namespace { + +/// GPU executors a search workspace with this configuration keeps on its device. +int searchWorkspaceExecutorCount(const SubstructSearchConfig& config) { + const int workers = config.workerThreads == -1 ? 4 : std::max(1, config.workerThreads); + return workers * resolveExecutorsPerRunner(config.executorsPerRunner, workers); +} + +} // namespace + +std::size_t estimateSubstructSearchWorkspaceBytes(const SubstructSearchConfig& config) { + const std::size_t batchSize = static_cast(std::max(1, config.batchSize)); + const std::size_t executors = static_cast(searchWorkspaceExecutorCount(config)); + const std::size_t overflowBuffers = config.algorithm == SubstructAlgorithm::GSI ? 2U : 1U; + const std::size_t overflowPerExecutor = + batchSize * overflowBuffers * static_cast(kOverflowEntriesPerBuffer) * sizeof(PartialMatch) * 3U / 2U; + const std::size_t auxiliaryPerExecutor = batchSize * 2048U + 8U * 1024U * 1024U; + return executors * (overflowPerExecutor + auxiliaryPerExecutor); +} + +std::size_t estimateRecursiveScratchBytes(const SubstructSearchConfig& config) { + // Recursive-pattern painting grows per-executor scratch for up to + // max(batchSize, 1024) blocks, each with two overflow buffers and a label + // matrix, by 1.5x (see RecursivePatternPreprocessor). + const std::size_t batchSize = static_cast(std::max(1, config.batchSize)); + const std::size_t paintBlocks = std::max(batchSize, 1024U); + const std::size_t perExecutor = paintBlocks * + (2U * static_cast(kOverflowEntriesPerBuffer) * sizeof(PartialMatch) + + kLabelMatrixWords * sizeof(std::uint32_t)) * + 3U / 2U; + return static_cast(searchWorkspaceExecutorCount(config)) * perExecutor; +} + } // namespace nvMolKit diff --git a/src/substruct/substruct_search.h b/src/substruct/substruct_search.h index bc92a44a..58bba7c2 100644 --- a/src/substruct/substruct_search.h +++ b/src/substruct/substruct_search.h @@ -18,6 +18,7 @@ #include +#include #include #include @@ -122,6 +123,15 @@ void hasSubstructMatch(const PersistentDeviceTargets& batch, /** Create reusable search state bound to one GPU for hasSubstructMatch(). */ std::shared_ptr makeSubstructSearchWorkspace(int deviceId); +/** Conservative device bytes one search workspace holds outside recursive queries. */ +std::size_t estimateSubstructSearchWorkspaceBytes(const SubstructSearchConfig& config); + +/** + * Additional device bytes a search workspace holds while a query with + * recursive SMARTS runs. It is released when the search returns. + */ +std::size_t estimateRecursiveScratchBytes(const SubstructSearchConfig& config); + } // namespace nvMolKit #endif // NVMOLKIT_SUBSTRUCTURE_SEARCH_H diff --git a/tests/test_substruct_library.cu b/tests/test_substruct_library.cu index 212fdd22..d37b5731 100644 --- a/tests/test_substruct_library.cu +++ b/tests/test_substruct_library.cu @@ -20,11 +20,14 @@ #include #include +#include #include +#include #include #include #include #include +#include #include #include @@ -203,6 +206,7 @@ TEST(SubstructLibraryState, FailedIncrementalFinalizeKeepsPreviousMoleculesAndCa library.addMol(*benzene); library.finalize(); + const std::size_t querySlots = library.maxConcurrentQueries(); library.addMol(*phenol); { @@ -212,6 +216,7 @@ TEST(SubstructLibraryState, FailedIncrementalFinalizeKeepsPreviousMoleculesAndCa EXPECT_EQ(library.size(), 1U); EXPECT_EQ(library.pendingSize(), 1U); + EXPECT_EQ(library.maxConcurrentQueries(), querySlots); EXPECT_EQ(library.getMatches(*aromaticCarbon), std::vector({0U})); library.finalize(); @@ -281,7 +286,7 @@ TEST(SubstructLibraryResults, ReusedWorkspaceDoesNotCarryMatchesBetweenQueries) } library.finalize(); - // Alternate a recursive and a plain query over the same reused workspace. + // Alternate a recursive and a plain query until every query slot has been reused. auto broadRecursive = queryFromSmarts("[$([#6])]"); auto narrowPlain = queryFromSmarts("[Cl-]"); ASSERT_NE(broadRecursive, nullptr); @@ -289,7 +294,7 @@ TEST(SubstructLibraryResults, ReusedWorkspaceDoesNotCarryMatchesBetweenQueries) const auto expectedBroad = rdkitMatchingIds(targets, *broadRecursive); const auto expectedNarrow = rdkitMatchingIds(targets, *narrowPlain); ASSERT_EQ(expectedNarrow, std::vector({2U})); - for (std::size_t round = 0; round < 3; ++round) { + for (std::size_t round = 0; round < 2 * library.maxConcurrentQueries() + 1; ++round) { EXPECT_EQ(library.getMatches(*broadRecursive), expectedBroad) << "round " << round; EXPECT_EQ(library.getMatches(*narrowPlain), expectedNarrow) << "round " << round; } @@ -497,6 +502,87 @@ TEST(SubstructLibraryMultiGpu, ShardsTargetsAndMergesEveryOperationInInsertionOr EXPECT_FALSE(library.hasMatch(*phosphorus)); } +TEST(SubstructLibraryConcurrency, MoreQueriesThanSlotsAllComplete) { + nvMolKit::SubstructLibrary library; + std::vector> targets; + for (const auto& smiles : {"CCO", "c1ccccc1", "CC(=O)O", "N", "OCCN"}) { + targets.push_back(molFromSmiles(smiles)); + ASSERT_NE(targets.back(), nullptr); + library.addMol(*targets.back()); + } + library.finalize(); + + auto recursive = queryFromSmarts("[$([#6]O)]"); + auto plain = queryFromSmarts("N"); + ASSERT_NE(recursive, nullptr); + ASSERT_NE(plain, nullptr); + const auto expectedRecursive = rdkitMatchingIds(targets, *recursive); + const auto expectedPlain = rdkitMatchingIds(targets, *plain); + + // More callers than query slots, half of them needing recursive scratch. + std::vector>> recursiveResults; + std::vector>> plainResults; + for (std::size_t caller = 0; caller < 2 * library.maxConcurrentQueries() + 2; ++caller) { + recursiveResults.push_back(std::async(std::launch::async, [&] { return library.getMatches(*recursive); })); + plainResults.push_back(std::async(std::launch::async, [&] { return library.getMatches(*plain); })); + } + for (auto& result : recursiveResults) { + EXPECT_EQ(result.get(), expectedRecursive); + } + for (auto& result : plainResults) { + EXPECT_EQ(result.get(), expectedPlain); + } +} + +TEST(SubstructLibraryConcurrency, QueriesSeeOnlyFinalizedMoleculesWhileMoreAreAdded) { + std::vector> targets; + for (const char* smiles : {"CCO", "c1ccccc1", "CCN", "OCCO", "c1ccncc1", "CC(=O)O", "CCCl", "Oc1ccccc1"}) { + targets.push_back(molFromSmiles(smiles)); + ASSERT_NE(targets.back(), nullptr); + } + auto oxygen = queryFromSmarts("[#8]"); + ASSERT_NE(oxygen, nullptr); + // A query sees the molecules of some finalize() call, never a partial set. + const std::vector finalizedSizes{2, 4, 6, 8}; + std::vector> finalizedResults; + for (const std::size_t size : finalizedSizes) { + std::vector> prefix; + for (std::size_t index = 0; index < size; ++index) { + prefix.push_back(std::make_unique(*targets[index])); + } + finalizedResults.push_back(rdkitMatchingIds(prefix, *oxygen)); + } + + nvMolKit::SubstructLibrary library; + library.addMols({targets[0].get(), targets[1].get()}); + library.finalize(); + + std::atomic done{false}; + std::vector> readers; + for (int reader = 0; reader < 3; ++reader) { + readers.push_back(std::async(std::launch::async, [&] { + int queries = 0; + while (!done.load()) { + const auto matches = library.getMatches(*oxygen); + EXPECT_NE(std::find(finalizedResults.begin(), finalizedResults.end(), matches), finalizedResults.end()); + ++queries; + } + return queries; + })); + } + for (std::size_t round = 1; round < finalizedSizes.size(); ++round) { + for (std::size_t index = finalizedSizes[round - 1]; index < finalizedSizes[round]; ++index) { + library.addMol(*targets[index]); + } + library.finalize(); + } + done.store(true); + for (auto& reader : readers) { + EXPECT_GT(reader.get(), 0); + } + EXPECT_EQ(library.getMatches(*oxygen), finalizedResults.back()); +} + TEST(SubstructLibraryStreams, SupportsFinalizeAndQueriesOnANondefaultStream) { nvMolKit::ScopedStream stream("SubstructLibraryStreams"); nvMolKit::SubstructLibrary library; From 757e78185d1bade489dd97e4ce464f23fd62a640 Mon Sep 17 00:00:00 2001 From: Kevin Boyd Date: Wed, 7 Oct 2026 17:10:28 -0400 Subject: [PATCH 2/4] Fix SubstructLibrary reader synchronization and reuse coverage --- tests/test_substruct_library.cu | 25 +++++++++++++++++++------ 1 file changed, 19 insertions(+), 6 deletions(-) diff --git a/tests/test_substruct_library.cu b/tests/test_substruct_library.cu index d37b5731..3c32a796 100644 --- a/tests/test_substruct_library.cu +++ b/tests/test_substruct_library.cu @@ -23,6 +23,7 @@ #include #include #include +#include #include #include #include @@ -286,7 +287,7 @@ TEST(SubstructLibraryResults, ReusedWorkspaceDoesNotCarryMatchesBetweenQueries) } library.finalize(); - // Alternate a recursive and a plain query until every query slot has been reused. + // Alternate recursive and plain queries to check sequential workspace reuse. auto broadRecursive = queryFromSmarts("[$([#6])]"); auto narrowPlain = queryFromSmarts("[Cl-]"); ASSERT_NE(broadRecursive, nullptr); @@ -294,7 +295,7 @@ TEST(SubstructLibraryResults, ReusedWorkspaceDoesNotCarryMatchesBetweenQueries) const auto expectedBroad = rdkitMatchingIds(targets, *broadRecursive); const auto expectedNarrow = rdkitMatchingIds(targets, *narrowPlain); ASSERT_EQ(expectedNarrow, std::vector({2U})); - for (std::size_t round = 0; round < 2 * library.maxConcurrentQueries() + 1; ++round) { + for (std::size_t round = 0; round < 3; ++round) { EXPECT_EQ(library.getMatches(*broadRecursive), expectedBroad) << "round " << round; EXPECT_EQ(library.getMatches(*narrowPlain), expectedNarrow) << "round " << round; } @@ -558,18 +559,30 @@ TEST(SubstructLibraryConcurrency, QueriesSeeOnlyFinalizedMoleculesWhileMoreAreAd library.finalize(); std::atomic done{false}; + std::latch readersReady{3}; std::vector> readers; for (int reader = 0; reader < 3; ++reader) { readers.push_back(std::async(std::launch::async, [&] { int queries = 0; - while (!done.load()) { - const auto matches = library.getMatches(*oxygen); - EXPECT_NE(std::find(finalizedResults.begin(), finalizedResults.end(), matches), finalizedResults.end()); - ++queries; + try { + while (!done.load()) { + const auto matches = library.getMatches(*oxygen); + EXPECT_NE(std::find(finalizedResults.begin(), finalizedResults.end(), matches), finalizedResults.end()); + if (++queries == 1) { + readersReady.count_down(); + } + } + } catch (...) { + // Let the writer finish so reader.get() can report a failed first query. + if (queries == 0) { + readersReady.count_down(); + } + throw; } return queries; })); } + readersReady.wait(); for (std::size_t round = 1; round < finalizedSizes.size(); ++round) { for (std::size_t index = finalizedSizes[round - 1]; index < finalizedSizes[round]; ++index) { library.addMol(*targets[index]); From c642a3e59fe533c6c611aee8ba5e3af0bbb5f1dd Mon Sep 17 00:00:00 2001 From: Kevin Boyd Date: Wed, 7 Oct 2026 21:00:22 -0400 Subject: [PATCH 3/4] Size SubstructLibrary query pools independently of OpenMP threads --- src/substruct/substruct_library.cpp | 5 +--- tests/test_substruct_library.cu | 45 +++++++++++++++++++++++++++++ 2 files changed, 46 insertions(+), 4 deletions(-) diff --git a/src/substruct/substruct_library.cpp b/src/substruct/substruct_library.cpp index 9d77938c..b9485eb1 100644 --- a/src/substruct/substruct_library.cpp +++ b/src/substruct/substruct_library.cpp @@ -450,10 +450,7 @@ class SubstructLibrary::Impl { } slots = std::min(slots, (availableBytes[gpu] - recursiveBytes[gpu]) / slotBytes[gpu]); } - // Each query searches every GPU from its own host thread. - const std::size_t hostThreads = static_cast(omp_get_max_threads()) / deviceIds_.size(); - slots = std::max(1, std::min(slots, hostThreads)); - std::size_t recursiveSlots = slots; + std::size_t recursiveSlots = slots; for (std::size_t gpu = 0; gpu < deviceIds_.size(); ++gpu) { recursiveSlots = std::min(recursiveSlots, (availableBytes[gpu] - slots * slotBytes[gpu]) / recursiveBytes[gpu]); } diff --git a/tests/test_substruct_library.cu b/tests/test_substruct_library.cu index 3c32a796..ee0f8739 100644 --- a/tests/test_substruct_library.cu +++ b/tests/test_substruct_library.cu @@ -18,6 +18,7 @@ #include #include #include +#include #include #include @@ -503,6 +504,50 @@ TEST(SubstructLibraryMultiGpu, ShardsTargetsAndMergesEveryOperationInInsertionOr EXPECT_FALSE(library.hasMatch(*phosphorus)); } +TEST(SubstructLibraryConcurrency, OpenMPThreadLimitDoesNotLimitQuerySlots) { + struct RestoreOpenMPThreads { + const int originalThreads = omp_get_max_threads(); + ~RestoreOpenMPThreads() { omp_set_num_threads(originalThreads); } + } restoreThreads; + + std::vector> targets; + for (const char* smiles : {"CCO", "c1ccccc1", "N", "OCCN"}) { + targets.push_back(molFromSmiles(smiles)); + ASSERT_NE(targets.back(), nullptr); + } + auto plain = queryFromSmarts("N"); + auto recursive = queryFromSmarts("[$([#6]O)]"); + ASSERT_NE(plain, nullptr); + ASSERT_NE(recursive, nullptr); + const auto expectedPlain = rdkitMatchingIds(targets, *plain); + const auto expectedRecursive = rdkitMatchingIds(targets, *recursive); + + for (const int threadCount : {1, 2}) { + omp_set_num_threads(threadCount); + for (const auto algorithm : {nvMolKit::SubstructAlgorithm::GSI, nvMolKit::SubstructAlgorithm::DFS}) { + SCOPED_TRACE(threadCount); + SCOPED_TRACE(static_cast(algorithm)); + nvMolKit::SubstructSearchConfig config; + config.algorithm = algorithm; + config.preprocessingThreads = 1; + config.workerThreads = 1; + config.executorsPerRunner = 1; + nvMolKit::SubstructLibrary library(config); + library.finalize(); + EXPECT_GT(library.maxConcurrentQueries(), 1U); + for (const auto& target : targets) { + library.addMol(*target); + } + library.finalize(); + EXPECT_GT(library.maxConcurrentQueries(), 1U); + auto plainResults = std::async(std::launch::async, [&] { return library.getMatches(*plain); }); + auto recursiveResults = std::async(std::launch::async, [&] { return library.getMatches(*recursive); }); + EXPECT_EQ(plainResults.get(), expectedPlain); + EXPECT_EQ(recursiveResults.get(), expectedRecursive); + } + } +} + TEST(SubstructLibraryConcurrency, MoreQueriesThanSlotsAllComplete) { nvMolKit::SubstructLibrary library; std::vector> targets; From 00f6d5292edbb80eae35be86a112d3283a1d9bf1 Mon Sep 17 00:00:00 2001 From: Kevin Boyd Date: Thu, 8 Oct 2026 09:18:29 -0400 Subject: [PATCH 4/4] Disable memory-dependent query concurrency regression test --- tests/test_substruct_library.cu | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tests/test_substruct_library.cu b/tests/test_substruct_library.cu index ee0f8739..27b2180f 100644 --- a/tests/test_substruct_library.cu +++ b/tests/test_substruct_library.cu @@ -504,7 +504,8 @@ TEST(SubstructLibraryMultiGpu, ShardsTargetsAndMergesEveryOperationInInsertionOr EXPECT_FALSE(library.hasMatch(*phosphorus)); } -TEST(SubstructLibraryConcurrency, OpenMPThreadLimitDoesNotLimitQuerySlots) { +// A shared GPU may validly fit only one query slot. This test needs a controlled memory budget to require more. +TEST(SubstructLibraryConcurrency, DISABLED_OpenMPThreadLimitDoesNotLimitQuerySlots) { struct RestoreOpenMPThreads { const int originalThreads = omp_get_max_threads(); ~RestoreOpenMPThreads() { omp_set_num_threads(originalThreads); }